Commit e2fbce20 e2fbce20b8c8dc21a07a4a10ed8771b06680339c by 沈秋雨

feat(etl): 新增 spider DB 歌手校验过滤无效录音

- 初始化阶段即校验 spider DB 歌手关联,过滤掉歌手数据缺失的录音
- 根据平台分别查询 QQ、酷狗和网易的有效歌手 key,过滤无效录音数据
- 统计并记录过滤掉的无效录音条数,日志输出详细过滤信息
- 文档更新新增 spider DB 歌手校验及相关数据统计说明
- 单元测试新增过滤无效歌手的情况,保证功能正确性
1 parent 5217dcac
# 数据统计与流转说明
> 统计时间:2026-07-13。records2 已按“不过滤 `pub_time`”的新规则重新初始化
> 统计时间:2026-07-14。records2 已按"不过滤 `pub_time`"的新规则重新初始化,并新增 spider DB 歌手校验
## 1. 主词曲录音关联表当前情况
......@@ -26,17 +26,20 @@ records1 当前共 **86,186** 条关联,且 **song_id** 也是 **86,186** 个
| 歌手非空后 | 208,368 | `hk_music_record.singer_name` |
| 歌曲链接非空后 | 206,201 | `hk_music_record.platform_index_url` |
| 音频链接非空后 | 206,201 | `storage_url``platform_play_url`,本次未减少 |
| 删除 records1 已有关联 | 120,949 | 删除 **85,252** 条相同 `(song_id, record_id)` 关联,输出最终 records2 |
| spider DB 歌手校验后 | 203,987 | 过滤 **2,214** 条歌手关联缺失的录音(QQ 为主) |
| 删除 records1 已有关联 | 118,848 | 删除 **85,139** 条相同 `(song_id, record_id)` 关联,输出最终 records2 |
# 数据规模汇总
| 数据集 | 数量 | 说明 |
| -------------------- | ----------: | ----------------------------- |
| records1 全量 | 86,186 | 每首歌一条主录音关联 |
| records2 全量 | 120,949 | 补充录音关联 |
| records2 全量(歌手校验后)| 118,848 | 补充录音关联(过滤 spider 歌手缺失)|
| 重复关联 | 0 | `(song_id, record_id)` 无重复 |
| **最终去重关联总数** | **207,135** | records1 + records2 |
| **最终去重关联总数** | **205,034** | records1 86,186 + records2 118,848 |
1. 已去除 records1 封面限制,已去除 records2 封面、歌词和 `pub_time` 限制,优先保留音频。records2 后续导入行为与 records1 保持一致,不对空 `pub_time` 做特殊转换。
2. 后续需给出各个环节数据缺失原因,判断是否可以人工补充。
2. spider DB 歌手校验:过滤 `media_tencent_singer_has_songs``singer_id = 0` 或无关联记录的录音。QQ 平台缺失最多(2,000 条),主要因为翻唱/二创版本的录音在爬虫库中歌手关联未正确建立。
3. 后续需给出各个环节数据缺失原因,判断是否可以人工补充。
......
......@@ -1409,11 +1409,54 @@ def initialize_yinyan_song_records(platforms: list[str], max_batches: int | None
log.info("Initialized yinyan_song_records candidates=%d, skipped_batches=%d", total, skipped_batches)
def _row_singer_key(row: dict) -> str | int:
"""返回一条 relation 用于匹配歌手的 key:QQ 用 mid(str),其他平台用 int id。"""
if row['platform'] == PLATFORM_QQ:
return row.get('platform_unique_key', '')
return int(row.get('platform_unique_key', 0))
def _fetch_valid_singer_keys(spider_conn, rows: list[dict]) -> set:
"""查询 spider DB,返回有有效歌手关联的 platform_unique_key 集合。"""
qq_mids = [r.get('platform_unique_key') for r in rows if r['platform'] == PLATFORM_QQ and r.get('platform_unique_key')]
kugou_ids = [int(r['platform_unique_key']) for r in rows if r['platform'] == PLATFORM_KUGOU and r.get('platform_unique_key')]
netease_ids = [int(r['platform_unique_key']) for r in rows if r['platform'] == PLATFORM_NETEASE and r.get('platform_unique_key')]
valid: set = set()
if qq_mids:
songs = fetch_qq_songs(spider_conn, qq_mids)
db_ids = [v['id'] for v in songs.values()]
singers = fetch_qq_singers(spider_conn, db_ids) if db_ids else {}
# singers 以 spider DB song_id 为 key,反查回 mid
id_to_mid = {v['id']: mid for mid, v in songs.items()}
for song_id in singers:
mid = id_to_mid.get(song_id)
if mid:
valid.add(mid)
if kugou_ids:
songs = fetch_kugou_songs(spider_conn, kugou_ids)
singers = fetch_kugou_singers(spider_conn, list(songs.keys()))
valid.update(int(k) for k in singers)
if netease_ids:
songs = fetch_netease_songs(spider_conn, netease_ids)
singers = fetch_netease_singers(spider_conn, list(songs.keys()))
valid.update(int(k) for k in singers)
return valid
def initialize_yinyan_song_records2(platforms: list[str], max_batches: int | None = None) -> None:
"""以 records1 全量 song_id 为范围,写入补充录音关联并去重。"""
"""以 records1 全量 song_id 为范围,写入补充录音关联并去重。
初始化阶段即校验 spider DB 歌手关联,过滤掉歌手数据缺失的录音。
"""
src_conn = get_source_conn()
spider_conn = get_spider_conn()
pg_conn = get_pg_conn()
total_inserted = 0
total_filtered_no_singer = 0
try:
with pg_conn.cursor() as pg_cur:
......@@ -1423,6 +1466,7 @@ def initialize_yinyan_song_records2(platforms: list[str], max_batches: int | Non
for index, start in enumerate(tqdm(range(0, len(song_ids), BATCH_SIZE), desc='init-yinyan-records2')):
# 长任务连接可能断开,每批次开始前刷新
src_conn = refresh_conn(src_conn, 'source')
spider_conn = refresh_conn(spider_conn, 'spider')
pg_conn = refresh_conn(pg_conn, 'pg')
if max_batches is not None and index >= max_batches:
......@@ -1435,6 +1479,7 @@ def initialize_yinyan_song_records2(platforms: list[str], max_batches: int | Non
'song_id': int(relation['source_song_id']),
'record_id': int(relation['record_id']),
'platform': str(relation['platform']),
'platform_unique_key': relation.get('platform_unique_key', ''),
}
for relation in relations
if str(relation['platform']) in platforms
......@@ -1442,6 +1487,20 @@ def initialize_yinyan_song_records2(platforms: list[str], max_batches: int | Non
if not rows:
continue
# ── 校验 spider DB 歌手关联,过滤无效录音 ──
valid_keys = _fetch_valid_singer_keys(spider_conn, rows)
before = len(rows)
rows = [
r for r in rows
if _row_singer_key(r) in valid_keys
]
filtered = before - len(rows)
total_filtered_no_singer += filtered
if filtered:
log.info("init-records2 batch %d: filtered %d rows without valid singers", index, filtered)
if not rows:
continue
with pg_conn.cursor() as pg_cur:
insert_yinyan_song_records2(pg_cur, rows)
pg_conn.commit()
......@@ -1455,9 +1514,11 @@ def initialize_yinyan_song_records2(platforms: list[str], max_batches: int | Non
close_all_pools()
log.info(
"Initialized yinyan_song_records2 candidates=%d, removed_existing_records1_relations=%d",
"Initialized yinyan_song_records2 candidates=%d, removed_existing_records1_relations=%d, "
"filtered_no_singer=%d",
total_inserted,
deleted,
total_filtered_no_singer,
)
......
......@@ -463,13 +463,15 @@ def test_initialize_yinyan_song_records2_writes_all_relations_then_deduplicates(
dedupe_calls = []
monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection())
monkeypatch.setattr(runner, 'get_spider_conn', lambda: _Connection())
monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn)
monkeypatch.setattr(runner, 'refresh_conn', lambda conn, name: conn)
monkeypatch.setattr(runner, 'fetch_yinyan_song_ids', lambda cur: [10])
monkeypatch.setattr(runner, 'fetch_all_song_record_relations', lambda conn, song_ids: [
{'source_song_id': 10, 'record_id': 100, 'platform': '1'},
{'source_song_id': 10, 'record_id': 200, 'platform': '2'},
{'source_song_id': 10, 'record_id': 100, 'platform': '1', 'platform_unique_key': 'mid100'},
{'source_song_id': 10, 'record_id': 200, 'platform': '2', 'platform_unique_key': '200'},
])
monkeypatch.setattr(runner, '_fetch_valid_singer_keys', lambda conn, rows: {'mid100', 200})
monkeypatch.setattr(runner, 'insert_yinyan_song_records2', lambda cur, rows: inserted.extend(rows))
monkeypatch.setattr(
runner,
......@@ -480,13 +482,42 @@ def test_initialize_yinyan_song_records2_writes_all_relations_then_deduplicates(
runner.initialize_yinyan_song_records2(['1', '2'])
assert inserted == [
{'song_id': 10, 'record_id': 100, 'platform': '1'},
{'song_id': 10, 'record_id': 200, 'platform': '2'},
{'song_id': 10, 'record_id': 100, 'platform': '1', 'platform_unique_key': 'mid100'},
{'song_id': 10, 'record_id': 200, 'platform': '2', 'platform_unique_key': '200'},
]
assert dedupe_calls == [True]
assert pg_conn.commits == 2
def test_initialize_yinyan_song_records2_filters_rows_without_valid_singers(monkeypatch):
pg_conn = _PgConnection()
inserted = []
monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection())
monkeypatch.setattr(runner, 'get_spider_conn', lambda: _Connection())
monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn)
monkeypatch.setattr(runner, 'refresh_conn', lambda conn, name: conn)
monkeypatch.setattr(runner, 'fetch_yinyan_song_ids', lambda cur: [10])
monkeypatch.setattr(runner, 'fetch_all_song_record_relations', lambda conn, song_ids: [
{'source_song_id': 10, 'record_id': 100, 'platform': '1', 'platform_unique_key': 'mid_good'},
{'source_song_id': 10, 'record_id': 200, 'platform': '1', 'platform_unique_key': 'mid_bad'},
{'source_song_id': 10, 'record_id': 300, 'platform': '2', 'platform_unique_key': '300'},
])
# mid_bad 没有歌手关联,300 也没有
monkeypatch.setattr(runner, '_fetch_valid_singer_keys', lambda conn, rows: {'mid_good'})
monkeypatch.setattr(runner, 'insert_yinyan_song_records2', lambda cur, rows: inserted.extend(rows))
monkeypatch.setattr(
runner, 'delete_yinyan_song_records2_existing_relations', lambda cur: 0,
)
runner.initialize_yinyan_song_records2(['1', '2'])
# 只有 mid_good 的录音被保留
assert inserted == [
{'song_id': 10, 'record_id': 100, 'platform': '1', 'platform_unique_key': 'mid_good'},
]
def test_backfill_yinyan_record_platforms_updates_missing_platform_rows(monkeypatch):
pg_conn = _PgConnection()
updated = []
......