perf(etl_to_crawler): 优化录音数据批量查询与写入性能
- 将批量查询录音数从1000调整为100,适配具体需求 - 新增按录音id批量查询录音函数,避免按歌曲id查询全部平台数据 - 修改runner流程中以录音id批量预取录音数据,减少重复查询 - 增加spider数据批量预取,避免逐条调用造成性能瓶颈 - 优化写入yinyan_song_records时采用批量更新语句,提高更新效率 - 新增测试覆盖fetch_platform_records_by_record_ids函数的查询准确性 - 调整测试用例以适配新增的批量查询函数及写入逻辑
Showing
7 changed files
with
124 additions
and
21 deletions
| ... | @@ -48,5 +48,5 @@ PLATFORM_KUGOU = '2' | ... | @@ -48,5 +48,5 @@ PLATFORM_KUGOU = '2' |
| 48 | PLATFORM_NETEASE = '4' | 48 | PLATFORM_NETEASE = '4' |
| 49 | PLATFORMS = [PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE] | 49 | PLATFORMS = [PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE] |
| 50 | 50 | ||
| 51 | BATCH_SIZE = 1000 | 51 | BATCH_SIZE = 100 |
| 52 | BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000')) | 52 | BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000')) | ... | ... |
| ... | @@ -54,6 +54,26 @@ ORDER BY sar.song_id, | ... | @@ -54,6 +54,26 @@ ORDER BY sar.song_id, |
| 54 | mr.id ASC | 54 | mr.id ASC |
| 55 | """ | 55 | """ |
| 56 | 56 | ||
| 57 | _PLATFORM_BY_RECORD_IDS_QUERY = """ | ||
| 58 | SELECT | ||
| 59 | sar.song_id AS source_song_id, | ||
| 60 | mr.id AS record_id, | ||
| 61 | mr.platform, | ||
| 62 | mr.platform_unique_key, | ||
| 63 | mr.platform_mid, | ||
| 64 | mr.album_audio_id, | ||
| 65 | sar.is_main_version, | ||
| 66 | mr.is_high, | ||
| 67 | mr.pub_time | ||
| 68 | FROM hk_song_and_record sar | ||
| 69 | JOIN hk_music_record mr ON mr.id = sar.record_id | ||
| 70 | WHERE mr.id IN ({placeholders}) | ||
| 71 | AND mr.platform IN ('1','2','4') | ||
| 72 | AND mr.platform_unique_key IS NOT NULL | ||
| 73 | AND mr.platform_unique_key != '' | ||
| 74 | AND mr.deleted = 0 | ||
| 75 | """ | ||
| 76 | |||
| 57 | # 同 _PLATFORM_QUERY,但不做 per-platform 去重,保留同平台全部录音供 singer fallback 遍历 | 77 | # 同 _PLATFORM_QUERY,但不做 per-platform 去重,保留同平台全部录音供 singer fallback 遍历 |
| 58 | _ALL_PLATFORM_RECORDS_QUERY = _PLATFORM_QUERY | 78 | _ALL_PLATFORM_RECORDS_QUERY = _PLATFORM_QUERY |
| 59 | 79 | ||
| ... | @@ -123,6 +143,27 @@ def fetch_all_platform_records(source_conn: pymysql.Connection, song_ids: list[i | ... | @@ -123,6 +143,27 @@ def fetch_all_platform_records(source_conn: pymysql.Connection, song_ids: list[i |
| 123 | return result | 143 | return result |
| 124 | 144 | ||
| 125 | 145 | ||
| 146 | def fetch_platform_records_by_record_ids(source_conn: pymysql.Connection, record_ids: list[int]) -> list[dict]: | ||
| 147 | """按状态表指定的录音 id 批量查询录音,不再按 song_id 拉取全平台候选。""" | ||
| 148 | if not record_ids: | ||
| 149 | return [] | ||
| 150 | unique_record_ids = list(dict.fromkeys(int(record_id) for record_id in record_ids)) | ||
| 151 | placeholders = ','.join(['%s'] * len(unique_record_ids)) | ||
| 152 | query = _PLATFORM_BY_RECORD_IDS_QUERY.format(placeholders=placeholders) | ||
| 153 | with source_conn.cursor() as cur: | ||
| 154 | cur.execute(query, unique_record_ids) | ||
| 155 | rows = cur.fetchall() | ||
| 156 | |||
| 157 | seen_record_ids: set = set() | ||
| 158 | result = [] | ||
| 159 | for row in rows: | ||
| 160 | rid = int(row['record_id']) | ||
| 161 | if rid not in seen_record_ids: | ||
| 162 | seen_record_ids.add(rid) | ||
| 163 | result.append(row) | ||
| 164 | return result | ||
| 165 | |||
| 166 | |||
| 126 | def fetch_record_platforms(source_conn: pymysql.Connection, record_ids: list[int]) -> dict[int, str]: | 167 | def fetch_record_platforms(source_conn: pymysql.Connection, record_ids: list[int]) -> dict[int, str]: |
| 127 | """按录音 id 查询平台代码,用于回填 yinyan_song_records.platform。""" | 168 | """按录音 id 查询平台代码,用于回填 yinyan_song_records.platform。""" |
| 128 | if not record_ids: | 169 | if not record_ids: | ... | ... |
| ... | @@ -10,6 +10,7 @@ from .reader import ( | ... | @@ -10,6 +10,7 @@ from .reader import ( |
| 10 | fetch_hk_songs_by_source_ids, | 10 | fetch_hk_songs_by_source_ids, |
| 11 | fetch_platform_records, | 11 | fetch_platform_records, |
| 12 | fetch_all_platform_records, | 12 | fetch_all_platform_records, |
| 13 | fetch_platform_records_by_record_ids, | ||
| 13 | fetch_record_platforms, | 14 | fetch_record_platforms, |
| 14 | select_primary_record, | 15 | select_primary_record, |
| 15 | ) | 16 | ) |
| ... | @@ -608,14 +609,40 @@ def run( | ... | @@ -608,14 +609,40 @@ def run( |
| 608 | break | 609 | break |
| 609 | 610 | ||
| 610 | song_ids = [int(r['song_id']) for r in pending_records] | 611 | song_ids = [int(r['song_id']) for r in pending_records] |
| 612 | record_ids = [int(r['record_id']) for r in pending_records] | ||
| 611 | hk_by_song = fetch_hk_songs_by_source_ids(hk_conn, song_ids) | 613 | hk_by_song = fetch_hk_songs_by_source_ids(hk_conn, song_ids) |
| 612 | all_platform_records = fetch_all_platform_records(src_conn, song_ids) | 614 | platform_records = fetch_platform_records_by_record_ids(src_conn, record_ids) |
| 613 | pr_by_state_key: dict[tuple[int, int, str], dict] = {} | 615 | pr_by_state_key: dict[tuple[int, int, str], dict] = {} |
| 614 | for pr in all_platform_records: | 616 | for pr in platform_records: |
| 615 | if pr['platform'] in platforms: | 617 | if pr['platform'] in platforms: |
| 616 | key = (int(pr['source_song_id']), int(pr['record_id']), pr['platform']) | 618 | key = (int(pr['source_song_id']), int(pr['record_id']), pr['platform']) |
| 617 | pr_by_state_key[key] = pr | 619 | pr_by_state_key[key] = pr |
| 618 | 620 | ||
| 621 | # ── 批量预取 spider 数据(整批一次 IN 查询,不逐首调用)────────── | ||
| 622 | qq_mids = list({pr['platform_unique_key'] for pr in platform_records if pr['platform'] == PLATFORM_QQ}) | ||
| 623 | kugou_ids = list({int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_KUGOU}) | ||
| 624 | netease_ids = list({int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_NETEASE}) | ||
| 625 | |||
| 626 | qq_songs_map = fetch_qq_songs(spider_conn, qq_mids) if qq_mids else {} | ||
| 627 | kugou_songs_map = fetch_kugou_songs(spider_conn, kugou_ids) if kugou_ids else {} | ||
| 628 | netease_songs_map = fetch_netease_songs(spider_conn, netease_ids) if netease_ids else {} | ||
| 629 | |||
| 630 | qq_db_song_ids = [v['id'] for v in qq_songs_map.values()] | ||
| 631 | qq_singers_map = fetch_qq_singers(spider_conn, qq_db_song_ids) if qq_db_song_ids else {} | ||
| 632 | kugou_singers_map = fetch_kugou_singers(spider_conn, kugou_ids) if kugou_ids else {} | ||
| 633 | netease_singers_map = fetch_netease_singers(spider_conn, netease_ids) if netease_ids else {} | ||
| 634 | |||
| 635 | songs_maps = { | ||
| 636 | PLATFORM_QQ: qq_songs_map, | ||
| 637 | PLATFORM_KUGOU: kugou_songs_map, | ||
| 638 | PLATFORM_NETEASE: netease_songs_map, | ||
| 639 | } | ||
| 640 | singers_maps = { | ||
| 641 | PLATFORM_QQ: qq_singers_map, | ||
| 642 | PLATFORM_KUGOU: kugou_singers_map, | ||
| 643 | PLATFORM_NETEASE: netease_singers_map, | ||
| 644 | } | ||
| 645 | |||
| 619 | with pg_conn.cursor() as pg_cur: | 646 | with pg_conn.cursor() as pg_cur: |
| 620 | pushed_count = 0 | 647 | pushed_count = 0 |
| 621 | for pending in pending_records: | 648 | for pending in pending_records: | ... | ... |
| ... | @@ -80,19 +80,22 @@ def upsert_yinyan_song_records(cur, records: list[dict]) -> None: | ... | @@ -80,19 +80,22 @@ def upsert_yinyan_song_records(cur, records: list[dict]) -> None: |
| 80 | """ | 80 | """ |
| 81 | if not records: | 81 | if not records: |
| 82 | return | 82 | return |
| 83 | cur.executemany( | 83 | values_sql = ', '.join(['(%s::bigint, %s::bigint, %s::varchar, %s::bigint)'] * len(records)) |
| 84 | """ | 84 | params = [] |
| 85 | UPDATE yinyan_song_records | 85 | for r in records: |
| 86 | SET platform = %s, | 86 | params.extend([r['song_id'], r['record_id'], r['platform'], r['platform_song_id']]) |
| 87 | platform_song_id = %s, | 87 | cur.execute( |
| 88 | record_id = %s, | 88 | f""" |
| 89 | UPDATE yinyan_song_records AS ysr | ||
| 90 | SET platform = v.platform, | ||
| 91 | platform_song_id = v.platform_song_id, | ||
| 92 | record_id = v.record_id, | ||
| 89 | is_yinyan_push = TRUE | 93 | is_yinyan_push = TRUE |
| 90 | WHERE song_id = %s AND is_yinyan_push = FALSE | 94 | FROM (VALUES {values_sql}) AS v(song_id, record_id, platform, platform_song_id) |
| 95 | WHERE ysr.song_id = v.song_id | ||
| 96 | AND ysr.is_yinyan_push = FALSE | ||
| 91 | """, | 97 | """, |
| 92 | [ | 98 | tuple(params), |
| 93 | (r['platform'], r['platform_song_id'], r['record_id'], r['song_id']) | ||
| 94 | for r in records | ||
| 95 | ], | ||
| 96 | ) | 99 | ) |
| 97 | 100 | ||
| 98 | 101 | ... | ... |
| 1 | from unittest.mock import MagicMock | 1 | from unittest.mock import MagicMock |
| 2 | 2 | ||
| 3 | from etl_to_crawler.reader import ( | 3 | from etl_to_crawler.reader import ( |
| 4 | fetch_platform_records_by_record_ids, | ||
| 4 | fetch_platform_records, | 5 | fetch_platform_records, |
| 5 | iter_hk_songs_batches, | 6 | iter_hk_songs_batches, |
| 6 | select_primary_record, | 7 | select_primary_record, |
| ... | @@ -59,6 +60,32 @@ def test_fetch_platform_records_keeps_one_record_per_song_and_platform(): | ... | @@ -59,6 +60,32 @@ def test_fetch_platform_records_keeps_one_record_per_song_and_platform(): |
| 59 | assert {row['record_id'] for row in result} == {100, 101, 102} | 60 | assert {row['record_id'] for row in result} == {100, 101, 102} |
| 60 | 61 | ||
| 61 | 62 | ||
| 63 | def test_fetch_platform_records_by_record_ids_queries_exact_records_once(): | ||
| 64 | conn = MagicMock() | ||
| 65 | conn.cursor.return_value = _make_cursor([ | ||
| 66 | { | ||
| 67 | 'source_song_id': 10, | ||
| 68 | 'record_id': 200, | ||
| 69 | 'platform': '2', | ||
| 70 | 'platform_unique_key': '200', | ||
| 71 | 'platform_mid': 'kg-hash', | ||
| 72 | 'album_audio_id': None, | ||
| 73 | 'is_main_version': 1, | ||
| 74 | 'is_high': 0, | ||
| 75 | 'pub_time': '2021-01-01', | ||
| 76 | }, | ||
| 77 | ]) | ||
| 78 | |||
| 79 | result = fetch_platform_records_by_record_ids(conn, [200, 200]) | ||
| 80 | |||
| 81 | cur = conn.cursor.return_value | ||
| 82 | sql, params = cur.execute.call_args[0] | ||
| 83 | assert 'WHERE mr.id IN (%s)' in sql | ||
| 84 | assert params == [200] | ||
| 85 | assert result[0]['record_id'] == 200 | ||
| 86 | assert result[0]['platform'] == '2' | ||
| 87 | |||
| 88 | |||
| 62 | def test_select_primary_record_prefers_main_version_then_is_high(): | 89 | def test_select_primary_record_prefers_main_version_then_is_high(): |
| 63 | rows = [ | 90 | rows = [ |
| 64 | { | 91 | { | ... | ... |
| ... | @@ -82,7 +82,8 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): | ... | @@ -82,7 +82,8 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): |
| 82 | }, | 82 | }, |
| 83 | ] | 83 | ] |
| 84 | monkeypatch.setattr(runner, 'fetch_platform_records', lambda conn, song_ids: platform_records) | 84 | monkeypatch.setattr(runner, 'fetch_platform_records', lambda conn, song_ids: platform_records) |
| 85 | monkeypatch.setattr(runner, 'fetch_all_platform_records', lambda conn, song_ids: platform_records) | 85 | monkeypatch.setattr(runner, 'fetch_all_platform_records', MagicMock()) |
| 86 | monkeypatch.setattr(runner, 'fetch_platform_records_by_record_ids', lambda conn, record_ids: platform_records) | ||
| 86 | monkeypatch.setattr(runner, '_PROCESSORS', processors) | 87 | monkeypatch.setattr(runner, '_PROCESSORS', processors) |
| 87 | monkeypatch.setattr(runner, 'upsert_yinyan_song_records', yinyan_writer) | 88 | monkeypatch.setattr(runner, 'upsert_yinyan_song_records', yinyan_writer) |
| 88 | 89 | ||
| ... | @@ -90,6 +91,7 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): | ... | @@ -90,6 +91,7 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): |
| 90 | 91 | ||
| 91 | processors['1'].assert_not_called() | 92 | processors['1'].assert_not_called() |
| 92 | processors['2'].assert_called_once() | 93 | processors['2'].assert_called_once() |
| 94 | runner.fetch_all_platform_records.assert_not_called() | ||
| 93 | yinyan_writer.assert_called_once_with(pg_conn.cur, [{ | 95 | yinyan_writer.assert_called_once_with(pg_conn.cur, [{ |
| 94 | 'song_id': 10, | 96 | 'song_id': 10, |
| 95 | 'record_id': 200, | 97 | 'record_id': 200, | ... | ... |
| ... | @@ -52,14 +52,17 @@ def test_upsert_yinyan_song_records_marks_existing_relation_as_pushed(): | ... | @@ -52,14 +52,17 @@ def test_upsert_yinyan_song_records_marks_existing_relation_as_pushed(): |
| 52 | {'song_id': 11, 'record_id': 101, 'platform': '2', 'platform_song_id': 1001}, | 52 | {'song_id': 11, 'record_id': 101, 'platform': '2', 'platform_song_id': 1001}, |
| 53 | ]) | 53 | ]) |
| 54 | 54 | ||
| 55 | sql, rows = cur.executemany.call_args[0] | 55 | sql, params = cur.execute.call_args[0] |
| 56 | assert 'UPDATE yinyan_song_records' in sql | 56 | assert 'UPDATE yinyan_song_records' in sql |
| 57 | assert 'platform = %s' in sql | 57 | assert 'UPDATE yinyan_song_records AS ysr' in sql |
| 58 | assert 'platform_song_id = %s' in sql | 58 | assert 'SET platform = v.platform' in sql |
| 59 | assert 'record_id = %s' in sql | 59 | assert 'platform_song_id = v.platform_song_id' in sql |
| 60 | assert 'record_id = v.record_id' in sql | ||
| 60 | assert 'is_yinyan_push = TRUE' in sql | 61 | assert 'is_yinyan_push = TRUE' in sql |
| 61 | assert 'WHERE song_id = %s AND is_yinyan_push = FALSE' in sql | 62 | assert 'FROM (VALUES (%s::bigint, %s::bigint, %s::varchar, %s::bigint), (%s::bigint, %s::bigint, %s::varchar, %s::bigint))' in sql |
| 62 | assert rows == [('1', 1000, 100, 10), ('2', 1001, 101, 11)] | 63 | assert 'ysr.song_id = v.song_id' in sql |
| 64 | assert 'ysr.is_yinyan_push = FALSE' in sql | ||
| 65 | assert params == (10, 100, '1', 1000, 11, 101, '2', 1001) | ||
| 63 | 66 | ||
| 64 | 67 | ||
| 65 | def test_insert_yinyan_song_records_initializes_unpushed_rows(): | 68 | def test_insert_yinyan_song_records_initializes_unpushed_rows(): | ... | ... |
-
Please register or sign in to post a comment