fix(etl_to_crawler): 解决同批次歌曲UUID冲突及关联关系处理问题
Showing
2 changed files
with
143 additions
and
6 deletions
| ... | @@ -442,7 +442,8 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing | ... | @@ -442,7 +442,8 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing |
| 442 | 'singers': singer_rows, | 442 | 'singers': singer_rows, |
| 443 | 'albums': album_rows, | 443 | 'albums': album_rows, |
| 444 | 'songs': [song_row], | 444 | 'songs': [song_row], |
| 445 | 'singer_songs': [(sg['singer_id'], song_uuid) for sg in singer_list], | 445 | # 关联阶段统一按平台歌曲 ID 解析数据库中的实际 UUID,兼容 songs 冲突跳过。 |
| 446 | 'singer_songs': [(sg['singer_id'], platform_song_id) for sg in singer_list], | ||
| 446 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], | 447 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], |
| 447 | } | 448 | } |
| 448 | 449 | ||
| ... | @@ -554,7 +555,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s | ... | @@ -554,7 +555,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s |
| 554 | 'singers': singer_rows, | 555 | 'singers': singer_rows, |
| 555 | 'albums': album_rows, | 556 | 'albums': album_rows, |
| 556 | 'songs': [song_row], | 557 | 'songs': [song_row], |
| 557 | 'singer_songs': [(sg['singer_id'], song_uuid) for sg in singer_list], | 558 | 'singer_songs': [(sg['singer_id'], song_id) for sg in singer_list], |
| 558 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], | 559 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], |
| 559 | } | 560 | } |
| 560 | 561 | ||
| ... | @@ -678,7 +679,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, | ... | @@ -678,7 +679,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, |
| 678 | 'singers': singer_rows, | 679 | 'singers': singer_rows, |
| 679 | 'albums': album_rows, | 680 | 'albums': album_rows, |
| 680 | 'songs': [song_row], | 681 | 'songs': [song_row], |
| 681 | 'singer_songs': [(sg['singer_id'], song_uuid) for sg in singer_list], | 682 | 'singer_songs': [(sg['singer_id'], song_id) for sg in singer_list], |
| 682 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], | 683 | 'singer_albums': [(sg['singer_id'], album_id) for sg in singer_list] if album_id else [], |
| 683 | } | 684 | } |
| 684 | 685 | ||
| ... | @@ -731,6 +732,57 @@ def _extend_platform_payload(target: dict, payload: dict) -> None: | ... | @@ -731,6 +732,57 @@ def _extend_platform_payload(target: dict, payload: dict) -> None: |
| 731 | target['singer_albums'].extend(payload.get('singer_albums', [])) | 732 | target['singer_albums'].extend(payload.get('singer_albums', [])) |
| 732 | 733 | ||
| 733 | 734 | ||
| 735 | _CRAWLER_SONG_TABLES = { | ||
| 736 | 'crawler_qqmusic_songs', | ||
| 737 | 'crawler_kugou_songs', | ||
| 738 | 'crawler_netease_songs', | ||
| 739 | } | ||
| 740 | |||
| 741 | |||
| 742 | def _dedupe_songs_by_platform_id(songs: list[dict]) -> list[dict]: | ||
| 743 | """同批相同平台歌曲只尝试插入一次,保留首次出现的数据。""" | ||
| 744 | seen = set() | ||
| 745 | deduped = [] | ||
| 746 | for song in songs: | ||
| 747 | platform_song_id = int(song['platform_song_id']) | ||
| 748 | if platform_song_id not in seen: | ||
| 749 | seen.add(platform_song_id) | ||
| 750 | deduped.append(song) | ||
| 751 | return deduped | ||
| 752 | |||
| 753 | |||
| 754 | def _resolve_singer_song_pairs( | ||
| 755 | cur, | ||
| 756 | songs_table: str, | ||
| 757 | platform_pairs: list[tuple[int, int]], | ||
| 758 | ) -> list[tuple[int, str]]: | ||
| 759 | """将 (singer_id, platform_song_id) 解析为关联表需要的实际歌曲 UUID。""" | ||
| 760 | if not platform_pairs: | ||
| 761 | return [] | ||
| 762 | if songs_table not in _CRAWLER_SONG_TABLES: | ||
| 763 | raise ValueError(f'unsupported crawler songs table: {songs_table}') | ||
| 764 | |||
| 765 | platform_song_ids = sorted({int(platform_song_id) for _, platform_song_id in platform_pairs}) | ||
| 766 | placeholders = ', '.join(['%s'] * len(platform_song_ids)) | ||
| 767 | cur.execute( | ||
| 768 | f'SELECT platform_song_id, id FROM {songs_table} ' | ||
| 769 | f'WHERE platform_song_id IN ({placeholders})', | ||
| 770 | tuple(platform_song_ids), | ||
| 771 | ) | ||
| 772 | actual_ids = {int(row[0]): str(row[1]) for row in cur.fetchall()} | ||
| 773 | missing = sorted(set(platform_song_ids) - actual_ids.keys()) | ||
| 774 | if missing: | ||
| 775 | raise RuntimeError( | ||
| 776 | f'cannot resolve actual song UUIDs from {songs_table}: platform_song_ids={missing}' | ||
| 777 | ) | ||
| 778 | |||
| 779 | # dict 保序去重,同一歌手与同一实际歌曲只写一次关系。 | ||
| 780 | return list(dict.fromkeys( | ||
| 781 | (int(singer_id), actual_ids[int(platform_song_id)]) | ||
| 782 | for singer_id, platform_song_id in platform_pairs | ||
| 783 | )) | ||
| 784 | |||
| 785 | |||
| 734 | def _write_import_payloads(pg_cur, payloads: list[dict]) -> list[dict]: | 786 | def _write_import_payloads(pg_cur, payloads: list[dict]) -> list[dict]: |
| 735 | if not payloads: | 787 | if not payloads: |
| 736 | return [] | 788 | return [] |
| ... | @@ -747,24 +799,36 @@ def _write_import_payloads(pg_cur, payloads: list[dict]) -> list[dict]: | ... | @@ -747,24 +799,36 @@ def _write_import_payloads(pg_cur, payloads: list[dict]) -> list[dict]: |
| 747 | imported.append(payload['result']) | 799 | imported.append(payload['result']) |
| 748 | 800 | ||
| 749 | qq = grouped[PLATFORM_QQ] | 801 | qq = grouped[PLATFORM_QQ] |
| 802 | qq['songs'] = _dedupe_songs_by_platform_id(qq['songs']) | ||
| 750 | upsert_qq_singers(pg_cur, qq['singers']) | 803 | upsert_qq_singers(pg_cur, qq['singers']) |
| 751 | upsert_qq_albums(pg_cur, qq['albums']) | 804 | upsert_qq_albums(pg_cur, qq['albums']) |
| 752 | upsert_qq_songs(pg_cur, qq['songs']) | 805 | upsert_qq_songs(pg_cur, qq['songs']) |
| 753 | upsert_qq_singer_songs(pg_cur, qq['singer_songs']) | 806 | qq_singer_songs = _resolve_singer_song_pairs( |
| 807 | pg_cur, 'crawler_qqmusic_songs', qq['singer_songs'], | ||
| 808 | ) | ||
| 809 | upsert_qq_singer_songs(pg_cur, qq_singer_songs) | ||
| 754 | upsert_qq_singer_albums(pg_cur, qq['singer_albums']) | 810 | upsert_qq_singer_albums(pg_cur, qq['singer_albums']) |
| 755 | 811 | ||
| 756 | kugou = grouped[PLATFORM_KUGOU] | 812 | kugou = grouped[PLATFORM_KUGOU] |
| 813 | kugou['songs'] = _dedupe_songs_by_platform_id(kugou['songs']) | ||
| 757 | upsert_kugou_singers(pg_cur, kugou['singers']) | 814 | upsert_kugou_singers(pg_cur, kugou['singers']) |
| 758 | upsert_kugou_albums(pg_cur, kugou['albums']) | 815 | upsert_kugou_albums(pg_cur, kugou['albums']) |
| 759 | upsert_kugou_songs(pg_cur, kugou['songs']) | 816 | upsert_kugou_songs(pg_cur, kugou['songs']) |
| 760 | upsert_kugou_singer_songs(pg_cur, kugou['singer_songs']) | 817 | kugou_singer_songs = _resolve_singer_song_pairs( |
| 818 | pg_cur, 'crawler_kugou_songs', kugou['singer_songs'], | ||
| 819 | ) | ||
| 820 | upsert_kugou_singer_songs(pg_cur, kugou_singer_songs) | ||
| 761 | upsert_kugou_singer_albums(pg_cur, kugou['singer_albums']) | 821 | upsert_kugou_singer_albums(pg_cur, kugou['singer_albums']) |
| 762 | 822 | ||
| 763 | netease = grouped[PLATFORM_NETEASE] | 823 | netease = grouped[PLATFORM_NETEASE] |
| 824 | netease['songs'] = _dedupe_songs_by_platform_id(netease['songs']) | ||
| 764 | upsert_netease_singers(pg_cur, netease['singers']) | 825 | upsert_netease_singers(pg_cur, netease['singers']) |
| 765 | upsert_netease_albums(pg_cur, netease['albums']) | 826 | upsert_netease_albums(pg_cur, netease['albums']) |
| 766 | upsert_netease_songs(pg_cur, netease['songs']) | 827 | upsert_netease_songs(pg_cur, netease['songs']) |
| 767 | upsert_netease_singer_songs(pg_cur, netease['singer_songs']) | 828 | netease_singer_songs = _resolve_singer_song_pairs( |
| 829 | pg_cur, 'crawler_netease_songs', netease['singer_songs'], | ||
| 830 | ) | ||
| 831 | upsert_netease_singer_songs(pg_cur, netease_singer_songs) | ||
| 768 | upsert_netease_singer_albums(pg_cur, netease['singer_albums']) | 832 | upsert_netease_singer_albums(pg_cur, netease['singer_albums']) |
| 769 | 833 | ||
| 770 | upsert_yinyan_song_records(pg_cur, yinyan_records) | 834 | upsert_yinyan_song_records(pg_cur, yinyan_records) | ... | ... |
| ... | @@ -241,6 +241,79 @@ def test_records2_selection_accepts_current_record_with_empty_cover_and_lyric(): | ... | @@ -241,6 +241,79 @@ def test_records2_selection_accepts_current_record_with_empty_cover_and_lyric(): |
| 241 | assert reason == 'ok' | 241 | assert reason == 'ok' |
| 242 | 242 | ||
| 243 | 243 | ||
| 244 | @pytest.mark.parametrize( | ||
| 245 | ('platform', 'songs_table', 'song_upsert_name', 'relation_upsert_name'), | ||
| 246 | [ | ||
| 247 | (runner.PLATFORM_QQ, 'crawler_qqmusic_songs', 'upsert_qq_songs', 'upsert_qq_singer_songs'), | ||
| 248 | (runner.PLATFORM_KUGOU, 'crawler_kugou_songs', 'upsert_kugou_songs', 'upsert_kugou_singer_songs'), | ||
| 249 | (runner.PLATFORM_NETEASE, 'crawler_netease_songs', 'upsert_netease_songs', 'upsert_netease_singer_songs'), | ||
| 250 | ], | ||
| 251 | ) | ||
| 252 | def test_write_import_payloads_resolves_existing_song_uuid_for_all_platforms( | ||
| 253 | monkeypatch, platform, songs_table, song_upsert_name, relation_upsert_name, | ||
| 254 | ): | ||
| 255 | cur = MagicMock() | ||
| 256 | cur.fetchall.return_value = [(200, 'existing-song-uuid')] | ||
| 257 | song_upsert = MagicMock() | ||
| 258 | relation_upsert = MagicMock() | ||
| 259 | yinyan_upsert = MagicMock() | ||
| 260 | monkeypatch.setattr(runner, song_upsert_name, song_upsert) | ||
| 261 | monkeypatch.setattr(runner, relation_upsert_name, relation_upsert) | ||
| 262 | monkeypatch.setattr(runner, 'upsert_yinyan_song_records', yinyan_upsert) | ||
| 263 | |||
| 264 | runner._write_import_payloads(cur, [{ | ||
| 265 | 'platform': platform, | ||
| 266 | 'platform_song_id': 200, | ||
| 267 | 'result': {'platform_song_id': 200}, | ||
| 268 | 'yinyan_record': { | ||
| 269 | 'song_id': 10, 'record_id': 20, 'platform': platform, 'platform_song_id': 200, | ||
| 270 | }, | ||
| 271 | 'singers': [], | ||
| 272 | 'albums': [], | ||
| 273 | 'songs': [ | ||
| 274 | {'song_uuid': 'temporary-uuid-1', 'platform_song_id': 200}, | ||
| 275 | {'song_uuid': 'temporary-uuid-2', 'platform_song_id': 200}, | ||
| 276 | ], | ||
| 277 | 'singer_songs': [(11, 200), (11, 200), (12, 200)], | ||
| 278 | 'singer_albums': [], | ||
| 279 | }]) | ||
| 280 | |||
| 281 | # 同批平台歌曲去重,且全部歌手关系使用数据库实际保留的歌曲 UUID。 | ||
| 282 | assert len(song_upsert.call_args.args[1]) == 1 | ||
| 283 | relation_upsert.assert_called_once_with( | ||
| 284 | cur, [(11, 'existing-song-uuid'), (12, 'existing-song-uuid')], | ||
| 285 | ) | ||
| 286 | sql, params = cur.execute.call_args.args | ||
| 287 | assert f'FROM {songs_table}' in sql | ||
| 288 | assert params == (200,) | ||
| 289 | yinyan_upsert.assert_called_once() | ||
| 290 | |||
| 291 | |||
| 292 | def test_write_import_payloads_does_not_mark_yinyan_success_when_song_uuid_is_missing(monkeypatch): | ||
| 293 | cur = MagicMock() | ||
| 294 | cur.fetchall.return_value = [] | ||
| 295 | yinyan_upsert = MagicMock() | ||
| 296 | monkeypatch.setattr(runner, 'upsert_yinyan_song_records', yinyan_upsert) | ||
| 297 | |||
| 298 | with pytest.raises(RuntimeError, match=r'platform_song_ids=\[200\]'): | ||
| 299 | runner._write_import_payloads(cur, [{ | ||
| 300 | 'platform': runner.PLATFORM_KUGOU, | ||
| 301 | 'platform_song_id': 200, | ||
| 302 | 'result': {'platform_song_id': 200}, | ||
| 303 | 'yinyan_record': { | ||
| 304 | 'song_id': 10, 'record_id': 20, 'platform': runner.PLATFORM_KUGOU, | ||
| 305 | 'platform_song_id': 200, | ||
| 306 | }, | ||
| 307 | 'singers': [], | ||
| 308 | 'albums': [], | ||
| 309 | 'songs': [{'song_uuid': 'temporary-uuid', 'platform_song_id': 200}], | ||
| 310 | 'singer_songs': [(11, 200)], | ||
| 311 | 'singer_albums': [], | ||
| 312 | }]) | ||
| 313 | |||
| 314 | yinyan_upsert.assert_not_called() | ||
| 315 | |||
| 316 | |||
| 244 | def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): | 317 | def test_run_imports_only_pending_yinyan_platform_record(monkeypatch): |
| 245 | pg_conn = _PgConnection() | 318 | pg_conn = _PgConnection() |
| 246 | prepare = MagicMock(return_value={ | 319 | prepare = MagicMock(return_value={ | ... | ... |
-
Please register or sign in to post a comment