Commit c41b71ef c41b71efbde226ed077cbbe5152a6da8bac78a93 by 沈秋雨

fix(runner): 调整导入资源必填规则并优化选择逻辑

- 修改 records2 歌曲基础字段要求,封面和歌词允许为空
- 根据导入类型动态设定必需资源,音频始终必需,歌词封面视情况而定
- 调整各平台录音选择函数,支持可选的必需资源参数
- 修改 _require_primary_assets 函数,支持传入必需资源参数校验
- 更新导入流程以使用动态资源需求,避免无必要的拒绝
- 添加单元测试覆盖不同导入表的资源要求和录音候选选择逻辑
- 代码注释及文档更新以反映新的资源校验规则
1 parent bf8d1723
...@@ -56,7 +56,7 @@ if YINYAN_IMPORT_TABLE not in _ALLOWED_YINYAN_IMPORT_TABLES: ...@@ -56,7 +56,7 @@ if YINYAN_IMPORT_TABLE not in _ALLOWED_YINYAN_IMPORT_TABLES:
56 f"got {YINYAN_IMPORT_TABLE!r}" 56 f"got {YINYAN_IMPORT_TABLE!r}"
57 ) 57 )
58 58
59 BATCH_SIZE = 1000 59 BATCH_SIZE = 10
60 BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000')) 60 BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000'))
61 HTTP_POOL_MAXSIZE = int(os.environ.get('HTTP_POOL_MAXSIZE', '128')) 61 HTTP_POOL_MAXSIZE = int(os.environ.get('HTTP_POOL_MAXSIZE', '128'))
62 HTTP_TRANSFER_RETRIES = int(os.environ.get('HTTP_TRANSFER_RETRIES', '3')) 62 HTTP_TRANSFER_RETRIES = int(os.environ.get('HTTP_TRANSFER_RETRIES', '3'))
......
...@@ -28,8 +28,7 @@ WHERE deleted = '0' ...@@ -28,8 +28,7 @@ WHERE deleted = '0'
28 ORDER BY id 28 ORDER BY id
29 """ 29 """
30 30
31 # 用于补全 yinyan_song_records2 的歌曲必须具备所有下游所需字段。 31 # records2 初始化不要求封面或歌词;正式导入时这两个字段允许为空。
32 # 当前 hk_songs_test 中的 lyrics_url 即歌曲链接字段。
33 _HK_SONGS_RECORDS2_QUERY = """ 32 _HK_SONGS_RECORDS2_QUERY = """
34 SELECT id, source_song_id 33 SELECT id, source_song_id
35 FROM hk_songs_test 34 FROM hk_songs_test
...@@ -39,7 +38,6 @@ WHERE deleted = '0' ...@@ -39,7 +38,6 @@ WHERE deleted = '0'
39 AND name IS NOT NULL AND TRIM(name) != '' 38 AND name IS NOT NULL AND TRIM(name) != ''
40 AND song_time IS NOT NULL AND song_time > 0 39 AND song_time IS NOT NULL AND song_time > 0
41 AND singer IS NOT NULL AND TRIM(singer) != '' 40 AND singer IS NOT NULL AND TRIM(singer) != ''
42 AND lyrics_url IS NOT NULL AND TRIM(lyrics_url) != ''
43 AND audio_url IS NOT NULL AND TRIM(audio_url) != '' 41 AND audio_url IS NOT NULL AND TRIM(audio_url) != ''
44 AND issue_time IS NOT NULL 42 AND issue_time IS NOT NULL
45 ORDER BY id 43 ORDER BY id
...@@ -150,7 +148,7 @@ def iter_hk_songs_records2_batches( ...@@ -150,7 +148,7 @@ def iter_hk_songs_records2_batches(
150 batch_size: int, 148 batch_size: int,
151 start_after_id: int = 0, 149 start_after_id: int = 0,
152 ) -> Iterator[list[dict]]: 150 ) -> Iterator[list[dict]]:
153 """遍历具备 records2 入库必填字段的歌曲。""" 151 """遍历具备 records2 基础入库字段的歌曲,封面和歌词可为空。"""
154 last_id = start_after_id 152 last_id = start_after_id
155 with conn.cursor() as cur: 153 with conn.cursor() as cur:
156 while True: 154 while True:
......
...@@ -4,7 +4,10 @@ import logging ...@@ -4,7 +4,10 @@ import logging
4 from concurrent.futures import ThreadPoolExecutor, as_completed 4 from concurrent.futures import ThreadPoolExecutor, as_completed
5 from tqdm import tqdm 5 from tqdm import tqdm
6 6
7 from .config import PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, BATCH_SIZE, BACKFILL_BATCH_SIZE, OSS_CONFIG 7 from .config import (
8 PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, BATCH_SIZE, BACKFILL_BATCH_SIZE,
9 OSS_CONFIG, YINYAN_IMPORT_TABLE,
10 )
8 from .connections import get_hk_songs_conn, get_source_conn, get_spider_conn, get_pg_conn, get_oss_bucket, refresh_conn, close_all_pools 11 from .connections import get_hk_songs_conn, get_source_conn, get_spider_conn, get_pg_conn, get_oss_bucket, refresh_conn, close_all_pools
9 from .reader import ( 12 from .reader import (
10 iter_hk_songs_batches, 13 iter_hk_songs_batches,
...@@ -121,14 +124,19 @@ def _safe_upload_lyric( ...@@ -121,14 +124,19 @@ def _safe_upload_lyric(
121 return '' 124 return ''
122 125
123 126
124 def _require_primary_assets(assets: dict) -> None: 127 def _required_import_assets() -> tuple[str, ...]:
125 """音频、歌曲封面、歌词均为入库必需资源,任一失败即拒绝整条记录。""" 128 """按导入状态表返回硬性资源要求;音频在两条支线中始终必需。"""
129 if YINYAN_IMPORT_TABLE == 'yinyan_song_records2':
130 return ('audio',)
131 return ('audio', 'lyric')
132
133
134 def _require_primary_assets(assets: dict, required: tuple[str, ...] | None = None) -> None:
135 """校验当前导入支线要求的资源,允许非必需资源以空值入库。"""
126 audio_url, _ = assets.get('audio', ('', '')) 136 audio_url, _ = assets.get('audio', ('', ''))
127 missing = [ 137 values = {'audio': audio_url, 'cover': assets.get('cover'), 'lyric': assets.get('lyric')}
128 name for name, value in ( 138 required = required or ('audio', 'cover', 'lyric')
129 ('audio', audio_url), ('cover', assets.get('cover')), ('lyric', assets.get('lyric')), 139 missing = [name for name in required if not values.get(name)]
130 ) if not value
131 ]
132 if missing: 140 if missing:
133 raise RequiredAssetTransferError(missing) 141 raise RequiredAssetTransferError(missing)
134 142
...@@ -206,14 +214,17 @@ def _select_qq_import_record( ...@@ -206,14 +214,17 @@ def _select_qq_import_record(
206 current_record: dict, 214 current_record: dict,
207 candidate_records: list[dict], 215 candidate_records: list[dict],
208 songs_by_mid: dict[str, dict], 216 songs_by_mid: dict[str, dict],
217 required_assets: tuple[str, ...] = ('lyric', 'cover'),
209 ) -> tuple[dict | None, str]: 218 ) -> tuple[dict | None, str]:
210 """为导入选择 QQ 录音;当前录音的歌词或封面不可用时,改选同源的合格候选。 219 """为导入选择 QQ 录音;当前录音的歌词或封面不可用时,改选同源的合格候选。
211 返回 (record, reason),record 为 None 时 reason 说明失败原因。""" 220 返回 (record, reason),record 为 None 时 reason 说明失败原因。"""
212 current_song = songs_by_mid.get(current_record['platform_unique_key']) 221 current_song = songs_by_mid.get(current_record['platform_unique_key'])
213 if not current_song: 222 if not current_song:
214 return None, 'spider data missing' 223 return None, 'spider data missing'
215 needs_lyric = not _has_usable_qq_lyric(current_song, hk_row) 224 needs_lyric = 'lyric' in required_assets and not _has_usable_qq_lyric(current_song, hk_row)
216 needs_cover = not _is_usable_qq_cover(_qq_cover_source(hk_row, current_song)) or not current_song.get('album_id') 225 needs_cover = 'cover' in required_assets and (
226 not _is_usable_qq_cover(_qq_cover_source(hk_row, current_song)) or not current_song.get('album_id')
227 )
217 if not needs_lyric and not needs_cover: 228 if not needs_lyric and not needs_cover:
218 return current_record, 'ok' 229 return current_record, 'ok'
219 230
...@@ -239,6 +250,7 @@ def _select_kugou_import_record( ...@@ -239,6 +250,7 @@ def _select_kugou_import_record(
239 current_record: dict, 250 current_record: dict,
240 candidate_records: list[dict], 251 candidate_records: list[dict],
241 songs_by_id: dict[int, dict], 252 songs_by_id: dict[int, dict],
253 required_assets: tuple[str, ...] = ('lyric', 'cover'),
242 ) -> tuple[dict | None, str]: 254 ) -> tuple[dict | None, str]:
243 """为导入选择酷狗录音;当前录音的歌词或封面不可用时,改选同源的合格候选。 255 """为导入选择酷狗录音;当前录音的歌词或封面不可用时,改选同源的合格候选。
244 返回 (record, reason),record 为 None 时 reason 说明失败原因。""" 256 返回 (record, reason),record 为 None 时 reason 说明失败原因。"""
...@@ -246,8 +258,8 @@ def _select_kugou_import_record( ...@@ -246,8 +258,8 @@ def _select_kugou_import_record(
246 current_song = songs_by_id.get(song_id) 258 current_song = songs_by_id.get(song_id)
247 if not current_song: 259 if not current_song:
248 return None, 'spider data missing' 260 return None, 'spider data missing'
249 needs_lyric = not _has_kugou_lyric(current_song, hk_row) 261 needs_lyric = 'lyric' in required_assets and not _has_kugou_lyric(current_song, hk_row)
250 needs_cover = not _has_kugou_cover(hk_row, current_song) 262 needs_cover = 'cover' in required_assets and not _has_kugou_cover(hk_row, current_song)
251 if not needs_lyric and not needs_cover: 263 if not needs_lyric and not needs_cover:
252 return current_record, 'ok' 264 return current_record, 'ok'
253 265
...@@ -274,6 +286,7 @@ def _select_netease_import_record( ...@@ -274,6 +286,7 @@ def _select_netease_import_record(
274 current_record: dict, 286 current_record: dict,
275 candidate_records: list[dict], 287 candidate_records: list[dict],
276 songs_by_id: dict[int, dict], 288 songs_by_id: dict[int, dict],
289 required_assets: tuple[str, ...] = ('lyric', 'cover'),
277 ) -> tuple[dict | None, str]: 290 ) -> tuple[dict | None, str]:
278 """为导入选择网易录音;当前录音的歌词或封面不可用时,改选同源的合格候选。 291 """为导入选择网易录音;当前录音的歌词或封面不可用时,改选同源的合格候选。
279 返回 (record, reason),record 为 None 时 reason 说明失败原因。""" 292 返回 (record, reason),record 为 None 时 reason 说明失败原因。"""
...@@ -281,8 +294,8 @@ def _select_netease_import_record( ...@@ -281,8 +294,8 @@ def _select_netease_import_record(
281 current_song = songs_by_id.get(song_id) 294 current_song = songs_by_id.get(song_id)
282 if not current_song: 295 if not current_song:
283 return None, 'spider data missing' 296 return None, 'spider data missing'
284 needs_lyric = not _has_netease_lyric(current_song, hk_row) 297 needs_lyric = 'lyric' in required_assets and not _has_netease_lyric(current_song, hk_row)
285 needs_cover = not _has_netease_cover(hk_row, current_song) 298 needs_cover = 'cover' in required_assets and not _has_netease_cover(hk_row, current_song)
286 if not needs_lyric and not needs_cover: 299 if not needs_lyric and not needs_cover:
287 return current_record, 'ok' 300 return current_record, 'ok'
288 301
...@@ -332,7 +345,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing ...@@ -332,7 +345,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
332 ) 345 )
333 ) 346 )
334 assets = _run_io_tasks(tasks) 347 assets = _run_io_tasks(tasks)
335 _require_primary_assets(assets) 348 _require_primary_assets(assets, _required_import_assets())
336 audio_url, audio_md5 = assets['audio'] 349 audio_url, audio_md5 = assets['audio']
337 cover_url = assets['cover'] 350 cover_url = assets['cover']
338 lyric_url = assets['lyric'] 351 lyric_url = assets['lyric']
...@@ -461,7 +474,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s ...@@ -461,7 +474,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s
461 ) 474 )
462 ) 475 )
463 assets = _run_io_tasks(tasks) 476 assets = _run_io_tasks(tasks)
464 _require_primary_assets(assets) 477 _require_primary_assets(assets, _required_import_assets())
465 audio_url, audio_md5 = assets['audio'] 478 audio_url, audio_md5 = assets['audio']
466 cover_url = assets['cover'] 479 cover_url = assets['cover']
467 lyric_url = assets['lyric'] 480 lyric_url = assets['lyric']
...@@ -573,7 +586,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, ...@@ -573,7 +586,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict,
573 ) 586 )
574 ) 587 )
575 assets = _run_io_tasks(tasks) 588 assets = _run_io_tasks(tasks)
576 _require_primary_assets(assets) 589 _require_primary_assets(assets, _required_import_assets())
577 audio_url, audio_md5 = assets['audio'] 590 audio_url, audio_md5 = assets['audio']
578 cover_url = assets['cover'] 591 cover_url = assets['cover']
579 lyric_url = assets['lyric'] 592 lyric_url = assets['lyric']
...@@ -1846,6 +1859,7 @@ def run( ...@@ -1846,6 +1859,7 @@ def run(
1846 1859
1847 prepare_inputs = [] 1860 prepare_inputs = []
1848 rejected_pending: list[dict] = [] 1861 rejected_pending: list[dict] = []
1862 required_assets = _required_import_assets()
1849 for pending in pending_records: 1863 for pending in pending_records:
1850 src_id = int(pending['song_id']) 1864 src_id = int(pending['song_id'])
1851 platform = str(pending['platform']) 1865 platform = str(pending['platform'])
...@@ -1866,15 +1880,15 @@ def run( ...@@ -1866,15 +1880,15 @@ def run(
1866 continue 1880 continue
1867 if platform == PLATFORM_QQ: 1881 if platform == PLATFORM_QQ:
1868 selected, reason = _select_qq_import_record( 1882 selected, reason = _select_qq_import_record(
1869 hk_row, pr, qq_records_by_song.get(src_id, []), qq_songs_map, 1883 hk_row, pr, qq_records_by_song.get(src_id, []), qq_songs_map, required_assets,
1870 ) 1884 )
1871 elif platform == PLATFORM_KUGOU: 1885 elif platform == PLATFORM_KUGOU:
1872 selected, reason = _select_kugou_import_record( 1886 selected, reason = _select_kugou_import_record(
1873 hk_row, pr, kugou_records_by_song.get(src_id, []), kugou_songs_map, 1887 hk_row, pr, kugou_records_by_song.get(src_id, []), kugou_songs_map, required_assets,
1874 ) 1888 )
1875 elif platform == PLATFORM_NETEASE: 1889 elif platform == PLATFORM_NETEASE:
1876 selected, reason = _select_netease_import_record( 1890 selected, reason = _select_netease_import_record(
1877 hk_row, pr, netease_records_by_song.get(src_id, []), netease_songs_map, 1891 hk_row, pr, netease_records_by_song.get(src_id, []), netease_songs_map, required_assets,
1878 ) 1892 )
1879 else: 1893 else:
1880 selected, reason = pr, 'ok' 1894 selected, reason = pr, 'ok'
......
...@@ -78,6 +78,33 @@ def test_required_asset_validation_rejects_partial_transfer(): ...@@ -78,6 +78,33 @@ def test_required_asset_validation_rejects_partial_transfer():
78 }) 78 })
79 79
80 80
81 def test_records_import_requires_audio_and_lyric_but_allows_empty_cover(monkeypatch):
82 monkeypatch.setattr(runner, 'YINYAN_IMPORT_TABLE', 'yinyan_song_records')
83
84 runner._require_primary_assets({
85 'audio': ('https://archive.example/audio.mp3', 'md5'),
86 'cover': '',
87 'lyric': 'https://archive.example/lyric.txt',
88 }, runner._required_import_assets())
89
90 with pytest.raises(runner.RequiredAssetTransferError, match='lyric'):
91 runner._require_primary_assets({
92 'audio': ('https://archive.example/audio.mp3', 'md5'),
93 'cover': '',
94 'lyric': '',
95 }, runner._required_import_assets())
96
97
98 def test_records2_import_allows_empty_cover_and_lyric(monkeypatch):
99 monkeypatch.setattr(runner, 'YINYAN_IMPORT_TABLE', 'yinyan_song_records2')
100
101 runner._require_primary_assets({
102 'audio': ('https://archive.example/audio.mp3', 'md5'),
103 'cover': '',
104 'lyric': '',
105 }, runner._required_import_assets())
106
107
81 def test_lyric_upload_prefers_hk_url_and_forces_oss_transfer(monkeypatch): 108 def test_lyric_upload_prefers_hk_url_and_forces_oss_transfer(monkeypatch):
82 download = MagicMock(return_value='[00:01.00]HK lyric') 109 download = MagicMock(return_value='[00:01.00]HK lyric')
83 monkeypatch.setattr(runner, 'download_text_url', download) 110 monkeypatch.setattr(runner, 'download_text_url', download)
...@@ -184,6 +211,36 @@ def test_select_qq_import_record_returns_none_when_no_candidate_satisfies_missin ...@@ -184,6 +211,36 @@ def test_select_qq_import_record_returns_none_when_no_candidate_satisfies_missin
184 assert 'lyric' in reason and 'cover' in reason 211 assert 'lyric' in reason and 'cover' in reason
185 212
186 213
214 def test_records_selection_does_not_reject_or_replace_for_missing_cover():
215 current = {'record_id': 10, 'platform_unique_key': 'current-mid'}
216
217 selected, reason = runner._select_qq_import_record(
218 {'cover_url': '', 'lyrics_url': 'https://example.com/hk-lyric.lrc'},
219 current,
220 [current],
221 {'current-mid': {'album_id': None, 'cover': '', 'lyric': ''}},
222 required_assets=('audio', 'lyric'),
223 )
224
225 assert selected == current
226 assert reason == 'ok'
227
228
229 def test_records2_selection_accepts_current_record_with_empty_cover_and_lyric():
230 current = {'record_id': 10, 'platform_unique_key': 'current-mid'}
231
232 selected, reason = runner._select_qq_import_record(
233 {'cover_url': '', 'lyrics_url': ''},
234 current,
235 [current],
236 {'current-mid': {'album_id': None, 'cover': '', 'lyric': ''}},
237 required_assets=('audio',),
238 )
239
240 assert selected == current
241 assert reason == 'ok'
242
243
187 def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): 244 def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
188 pg_conn = _PgConnection() 245 pg_conn = _PgConnection()
189 prepare = MagicMock(return_value={ 246 prepare = MagicMock(return_value={
...@@ -206,6 +263,7 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): ...@@ -206,6 +263,7 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
206 monkeypatch.setattr(runner, 'get_spider_conn', lambda: _Connection()) 263 monkeypatch.setattr(runner, 'get_spider_conn', lambda: _Connection())
207 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn) 264 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn)
208 monkeypatch.setattr(runner, 'get_oss_bucket', lambda: object()) 265 monkeypatch.setattr(runner, 'get_oss_bucket', lambda: object())
266 monkeypatch.setattr(runner, 'refresh_conn', lambda conn, name: conn)
209 monkeypatch.setattr(runner, 'fetch_pending_yinyan_song_records', lambda cur, batch_size: [ 267 monkeypatch.setattr(runner, 'fetch_pending_yinyan_song_records', lambda cur, batch_size: [
210 {'song_id': 10, 'record_id': 200, 'platform': '2'}, 268 {'song_id': 10, 'record_id': 200, 'platform': '2'},
211 ] if pg_conn.commits == 0 else []) 269 ] if pg_conn.commits == 0 else [])
...@@ -262,6 +320,7 @@ def test_initialize_yinyan_song_records_inserts_primary_records(monkeypatch): ...@@ -262,6 +320,7 @@ def test_initialize_yinyan_song_records_inserts_primary_records(monkeypatch):
262 monkeypatch.setattr(runner, 'get_hk_songs_conn', lambda: _Connection()) 320 monkeypatch.setattr(runner, 'get_hk_songs_conn', lambda: _Connection())
263 monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection()) 321 monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection())
264 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn) 322 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn)
323 monkeypatch.setattr(runner, 'refresh_conn', lambda conn, name: conn)
265 monkeypatch.setattr(runner, 'fetch_existing_yinyan_song_ids', lambda cur: set()) 324 monkeypatch.setattr(runner, 'fetch_existing_yinyan_song_ids', lambda cur: set())
266 325
267 monkeypatch.setattr(runner, 'iter_hk_songs_batches', lambda conn, batch_size: [[ 326 monkeypatch.setattr(runner, 'iter_hk_songs_batches', lambda conn, batch_size: [[
...@@ -335,6 +394,7 @@ def test_backfill_yinyan_record_platforms_updates_missing_platform_rows(monkeypa ...@@ -335,6 +394,7 @@ def test_backfill_yinyan_record_platforms_updates_missing_platform_rows(monkeypa
335 394
336 monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection()) 395 monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection())
337 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn) 396 monkeypatch.setattr(runner, 'get_pg_conn', lambda: pg_conn)
397 monkeypatch.setattr(runner, 'refresh_conn', lambda conn, name: conn)
338 monkeypatch.setattr(runner, 'fetch_yinyan_records_missing_platform', lambda cur, batch_size: [ 398 monkeypatch.setattr(runner, 'fetch_yinyan_records_missing_platform', lambda cur, batch_size: [
339 {'song_id': 10, 'record_id': 100}, 399 {'song_id': 10, 'record_id': 100},
340 {'song_id': 11, 'record_id': 101}, 400 {'song_id': 11, 'record_id': 101},
......