Commit ac91b0fe ac91b0fe0e1d6623f8af9968a5a56a1dab40342e by 沈秋雨

feat(crawler): 为酷狗数据新增provider和源数据字段支持

- 新增 PROVIDER_KUGOU 常量用于标识数据提供者
- 增加 _source_json 方法以支持复杂数据的JSON序列化
- 在处理酷狗歌手数据时添加provider_name和crawler_source_data字段
- 在处理酷狗专辑和歌曲数据时添加对应的provider_name和crawler_source_data字段
- netease歌曲数据处理新增album_json字段以存储专辑数据的JSON表示
- 写入数据库时酷狗歌手、专辑、歌曲表增加provider_name和crawler_source_data列及其对应数据
- 网易歌曲写入增加album_json字段支持
- 新增单元测试覆盖酷狗及网易歌曲处理与数据库写入的新增字段和逻辑
1 parent e04906eb
...@@ -26,6 +26,7 @@ from .lyric import strip_timestamps ...@@ -26,6 +26,7 @@ from .lyric import strip_timestamps
26 26
27 logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s') 27 logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
28 log = logging.getLogger(__name__) 28 log = logging.getLogger(__name__)
29 PROVIDER_KUGOU = 'kugou'
29 30
30 31
31 def _safe_transfer(url, oss_key, bucket, base_url): 32 def _safe_transfer(url, oss_key, bucket, base_url):
...@@ -36,6 +37,14 @@ def _safe_transfer(url, oss_key, bucket, base_url): ...@@ -36,6 +37,14 @@ def _safe_transfer(url, oss_key, bucket, base_url):
36 return url # 失败时保留原 URL,不阻断流程 37 return url # 失败时保留原 URL,不阻断流程
37 38
38 39
40 def _json_default(value):
41 return str(value)
42
43
44 def _source_json(data: dict) -> str:
45 return json.dumps(data, ensure_ascii=False, default=_json_default)
46
47
39 def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url): 48 def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
40 mid = pr['platform_unique_key'] 49 mid = pr['platform_unique_key']
41 songs_map = fetch_qq_songs(spider_conn, [mid]) 50 songs_map = fetch_qq_songs(spider_conn, [mid])
...@@ -165,7 +174,13 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url ...@@ -165,7 +174,13 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
165 build_oss_key('kugou', 'singer', str(sg['singer_id']) + '.jpg'), 174 build_oss_key('kugou', 'singer', str(sg['singer_id']) + '.jpg'),
166 bucket, base_url 175 bucket, base_url
167 ) 176 )
168 singer_rows.append({**sg, 'id': sg['singer_id'], 'avatar': avatar}) 177 singer_rows.append({
178 **sg,
179 'id': sg['singer_id'],
180 'avatar': avatar,
181 'provider_name': PROVIDER_KUGOU,
182 'crawler_source_data': _source_json(sg),
183 })
169 upsert_kugou_singers(pg_cur, singer_rows) 184 upsert_kugou_singers(pg_cur, singer_rows)
170 185
171 singers_json = json.dumps([{ 186 singers_json = json.dumps([{
...@@ -183,6 +198,18 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url ...@@ -183,6 +198,18 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
183 'type': sp.get('album_type') or '', 'company_id': sp.get('company_id') or 0, 198 'type': sp.get('album_type') or '', 'company_id': sp.get('company_id') or 0,
184 'company': sp.get('company') or '', 'is_owner': sp.get('is_owner') or 0, 199 'company': sp.get('company') or '', 'is_owner': sp.get('is_owner') or 0,
185 'published_at': sp.get('album_published_at'), 200 'published_at': sp.get('album_published_at'),
201 'provider_name': PROVIDER_KUGOU,
202 'crawler_source_data': _source_json({
203 'id': sp.get('album_id'),
204 'cover': sp.get('album_cover'),
205 'title': sp.get('album_title'),
206 'intro': sp.get('album_intro'),
207 'type': sp.get('album_type'),
208 'company_id': sp.get('company_id'),
209 'company': sp.get('company'),
210 'is_owner': sp.get('is_owner'),
211 'published_at': sp.get('album_published_at'),
212 }),
186 }]) 213 }])
187 214
188 song_uuid = str(uuid.uuid4()) 215 song_uuid = str(uuid.uuid4())
...@@ -204,6 +231,8 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url ...@@ -204,6 +231,8 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
204 'platform_index_url': sp.get('platform_index_url') or f'http://www.kugou.com/song/#hash={sp.get("hid") or pr.get("platform_mid", "")}', 231 'platform_index_url': sp.get('platform_index_url') or f'http://www.kugou.com/song/#hash={sp.get("hid") or pr.get("platform_mid", "")}',
205 'published_at': sp.get('published_at') or hk_row.get('issue_time'), 232 'published_at': sp.get('published_at') or hk_row.get('issue_time'),
206 'singers_json': singers_json, 233 'singers_json': singers_json,
234 'provider_name': PROVIDER_KUGOU,
235 'crawler_source_data': _source_json(sp),
207 }]) 236 }])
208 # 查询实际 UUID(ON CONFLICT DO NOTHING 时使用已有 UUID) 237 # 查询实际 UUID(ON CONFLICT DO NOTHING 时使用已有 UUID)
209 pg_cur.execute('SELECT id FROM crawler_kugou_songs WHERE platform_song_id = %s', (song_id,)) 238 pg_cur.execute('SELECT id FROM crawler_kugou_songs WHERE platform_song_id = %s', (song_id,))
...@@ -262,8 +291,21 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u ...@@ -262,8 +291,21 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u
262 } for sg in singer_list], ensure_ascii=False) 291 } for sg in singer_list], ensure_ascii=False)
263 292
264 album_id = None 293 album_id = None
294 album_json = None
265 if sp.get('album_id'): 295 if sp.get('album_id'):
266 album_id = sp['album_id'] 296 album_id = sp['album_id']
297 album_payload = {
298 'id': sp['album_id'],
299 'cover': album_cover,
300 'title': sp.get('album_title') or '',
301 'intro': sp.get('album_intro'),
302 'type': sp.get('album_type') or '',
303 'company_id': sp.get('company_id') or 0,
304 'company': sp.get('company') or '',
305 'is_owner': sp.get('is_owner') or 0,
306 'published_at': str(sp.get('album_published_at')) if sp.get('album_published_at') else None,
307 }
308 album_json = json.dumps(album_payload, ensure_ascii=False)
267 upsert_netease_albums(pg_cur, [{ 309 upsert_netease_albums(pg_cur, [{
268 'id': sp['album_id'], 'cover': album_cover, 310 'id': sp['album_id'], 'cover': album_cover,
269 'title': sp.get('album_title') or '', 'intro': sp.get('album_intro'), 311 'title': sp.get('album_title') or '', 'intro': sp.get('album_intro'),
...@@ -277,6 +319,7 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u ...@@ -277,6 +319,7 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u
277 'song_uuid': song_uuid, 319 'song_uuid': song_uuid,
278 'platform_song_id': song_id, 320 'platform_song_id': song_id,
279 'album_id': album_id, 321 'album_id': album_id,
322 'album_json': album_json,
280 'cover': cover_url, 323 'cover': cover_url,
281 'title': sp.get('title', hk_row['name']), 324 'title': sp.get('title', hk_row['name']),
282 'name': hk_row['name'], 325 'name': hk_row['name'],
......
...@@ -108,14 +108,17 @@ def upsert_kugou_singers(cur, singers: list[dict]) -> None: ...@@ -108,14 +108,17 @@ def upsert_kugou_singers(cur, singers: list[dict]) -> None:
108 return 108 return
109 sql = """ 109 sql = """
110 INSERT INTO crawler_kugou_singers 110 INSERT INTO crawler_kugou_singers
111 (id, name, avatar, sex, area, "index", intro, home_url, created_at, updated_at) 111 (id, name, avatar, sex, area, "index", intro, home_url,
112 VALUES (%s, %s, %s, %s::kugou_singer_sex, %s::kugou_singer_area, %s::kugou_singer_index, %s, %s, NOW(), NOW()) 112 created_at, updated_at, provider_name, crawler_source_data)
113 VALUES (%s, %s, %s, %s::kugou_singer_sex, %s::kugou_singer_area, %s::kugou_singer_index, %s, %s,
114 NOW(), NOW(), %s, %s::json)
113 ON CONFLICT (id) DO NOTHING 115 ON CONFLICT (id) DO NOTHING
114 """ 116 """
115 rows = [( 117 rows = [(
116 s['id'], s['name'], s.get('avatar', ''), 118 s['id'], s['name'], s.get('avatar', ''),
117 s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#', 119 s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
118 s.get('intro'), s.get('home_url', ''), 120 s.get('intro'), s.get('home_url', ''),
121 s.get('provider_name'), s.get('crawler_source_data'),
119 ) for s in singers] 122 ) for s in singers]
120 cur.executemany(sql, rows) 123 cur.executemany(sql, rows)
121 124
...@@ -125,14 +128,16 @@ def upsert_kugou_albums(cur, albums: list[dict]) -> None: ...@@ -125,14 +128,16 @@ def upsert_kugou_albums(cur, albums: list[dict]) -> None:
125 return 128 return
126 sql = """ 129 sql = """
127 INSERT INTO crawler_kugou_albums 130 INSERT INTO crawler_kugou_albums
128 (id, cover, title, intro, type, company_id, company, is_owner, published_at, created_at, updated_at) 131 (id, cover, title, intro, type, company_id, company, is_owner, published_at,
129 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW()) 132 created_at, updated_at, provider_name, crawler_source_data)
133 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
130 ON CONFLICT (id) DO NOTHING 134 ON CONFLICT (id) DO NOTHING
131 """ 135 """
132 rows = [( 136 rows = [(
133 a['id'], a.get('cover', ''), a.get('title', ''), a.get('intro'), 137 a['id'], a.get('cover', ''), a.get('title', ''), a.get('intro'),
134 a.get('type', ''), a.get('company_id', 0) or 0, 138 a.get('type', ''), a.get('company_id', 0) or 0,
135 a.get('company', ''), a.get('is_owner', 0) or 0, a.get('published_at'), 139 a.get('company', ''), a.get('is_owner', 0) or 0, a.get('published_at'),
140 a.get('provider_name'), a.get('crawler_source_data'),
136 ) for a in albums] 141 ) for a in albums]
137 cur.executemany(sql, rows) 142 cur.executemany(sql, rows)
138 143
...@@ -144,8 +149,9 @@ def upsert_kugou_songs(cur, songs: list[dict]) -> None: ...@@ -144,8 +149,9 @@ def upsert_kugou_songs(cur, songs: list[dict]) -> None:
144 INSERT INTO crawler_kugou_songs 149 INSERT INTO crawler_kugou_songs
145 (id, platform_song_id, hash, album_audio_id, album_id, cover, title, name, duration, 150 (id, platform_song_id, hash, album_audio_id, album_id, cover, title, name, duration,
146 lyric, composer_name, lyricist_name, url, lyric_url, 151 lyric, composer_name, lyricist_name, url, lyric_url,
147 platform_index_url, published_at, singers, status, created_at, updated_at) 152 platform_index_url, published_at, singers, status, created_at, updated_at,
148 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW()) 153 provider_name, crawler_source_data)
154 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW(), %s, %s::json)
149 ON CONFLICT (platform_song_id) DO NOTHING 155 ON CONFLICT (platform_song_id) DO NOTHING
150 """ 156 """
151 rows = [( 157 rows = [(
...@@ -155,6 +161,7 @@ def upsert_kugou_songs(cur, songs: list[dict]) -> None: ...@@ -155,6 +161,7 @@ def upsert_kugou_songs(cur, songs: list[dict]) -> None:
155 s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'), 161 s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'),
156 s.get('url', ''), s.get('lyric_url'), 162 s.get('url', ''), s.get('lyric_url'),
157 s.get('platform_index_url'), s.get('published_at'), s.get('singers_json', '[]'), 163 s.get('platform_index_url'), s.get('published_at'), s.get('singers_json', '[]'),
164 s.get('provider_name'), s.get('crawler_source_data'),
158 ) for s in songs] 165 ) for s in songs]
159 cur.executemany(sql, rows) 166 cur.executemany(sql, rows)
160 167
...@@ -224,8 +231,8 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None: ...@@ -224,8 +231,8 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None:
224 INSERT INTO crawler_netease_songs 231 INSERT INTO crawler_netease_songs
225 (id, platform_song_id, album_id, cover, title, name, duration, 232 (id, platform_song_id, album_id, cover, title, name, duration,
226 lyric, composer_name, lyricist_name, url, lyric_url, 233 lyric, composer_name, lyricist_name, url, lyric_url,
227 platform_index_url, published_at, singers, status, created_at, updated_at) 234 platform_index_url, published_at, album, singers, status, created_at, updated_at)
228 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW()) 235 VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW())
229 ON CONFLICT (platform_song_id) DO NOTHING 236 ON CONFLICT (platform_song_id) DO NOTHING
230 """ 237 """
231 rows = [( 238 rows = [(
...@@ -233,7 +240,7 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None: ...@@ -233,7 +240,7 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None:
233 s.get('cover', ''), s.get('title', ''), s.get('name', ''), s.get('duration', 0) or 0, 240 s.get('cover', ''), s.get('title', ''), s.get('name', ''), s.get('duration', 0) or 0,
234 s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'), 241 s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'),
235 s.get('url', ''), s.get('lyric_url'), 242 s.get('url', ''), s.get('lyric_url'),
236 s.get('platform_index_url'), s.get('published_at'), s.get('singers_json', '[]'), 243 s.get('platform_index_url'), s.get('published_at'), s.get('album_json'), s.get('singers_json', '[]'),
237 ) for s in songs] 244 ) for s in songs]
238 cur.executemany(sql, rows) 245 cur.executemany(sql, rows)
239 246
......
...@@ -116,3 +116,61 @@ def test_run_resume_starts_after_saved_id_and_updates_state_after_commit(monkeyp ...@@ -116,3 +116,61 @@ def test_run_resume_starts_after_saved_id_and_updates_state_after_commit(monkeyp
116 assert seen_start_ids == [40] 116 assert seen_start_ids == [40]
117 assert pg_conn.commits == 1 117 assert pg_conn.commits == 1
118 assert '"last_hk_songs_id": 60' in state_file.read_text(encoding='utf-8') 118 assert '"last_hk_songs_id": 60' in state_file.read_text(encoding='utf-8')
119
120
121 def test_process_netease_builds_album_json_for_song_insert(monkeypatch):
122 pg_cur = MagicMock()
123 pg_cur.fetchone.return_value = ('song-uuid',)
124 inserted_songs = []
125
126 monkeypatch.setattr(runner, 'fetch_netease_songs', lambda conn, song_ids: {
127 300: {
128 'id': 300,
129 'album_id': 20,
130 'album_cover': 'https://example.com/album.jpg',
131 'album_title': '专辑',
132 'album_intro': '简介',
133 'album_type': '专辑类型',
134 'company_id': 7,
135 'company': '唱片公司',
136 'is_owner': 1,
137 'album_published_at': '2020-01-01',
138 'cover': 'https://example.com/cover.jpg',
139 'title': '录音标题',
140 'duration': 180,
141 'lyric': '[00:01.00]歌词',
142 'composer_name': '曲作者',
143 'lyricist_name': '词作者',
144 'platform_index_url': None,
145 'published_at': '2020-01-02',
146 },
147 })
148 monkeypatch.setattr(runner, 'fetch_netease_singers', lambda conn, song_ids: {})
149 monkeypatch.setattr(runner, '_safe_transfer', lambda url, oss_key, bucket, base_url: url)
150 monkeypatch.setattr(runner, 'upsert_netease_singers', lambda cur, singers: None)
151 monkeypatch.setattr(runner, 'upsert_netease_albums', lambda cur, albums: None)
152 monkeypatch.setattr(runner, 'upsert_netease_songs', lambda cur, songs: inserted_songs.extend(songs))
153 monkeypatch.setattr(runner, 'upsert_netease_singer_songs', lambda cur, pairs: None)
154 monkeypatch.setattr(runner, 'upsert_netease_singer_albums', lambda cur, pairs: None)
155
156 runner._process_netease(
157 {
158 'name': '词曲名',
159 'audio_url': 'https://example.com/audio.mp3',
160 'lyrics_url': 'https://example.com/lyric.lrc',
161 'cover_url': '',
162 'composer': '词曲曲作者',
163 'lyricist': '词曲词作者',
164 'issue_time': '2019-01-01',
165 'song_time': 120,
166 },
167 {'platform_unique_key': '300'},
168 spider_conn=object(),
169 pg_cur=pg_cur,
170 bucket=object(),
171 base_url='https://bucket.example.com',
172 )
173
174 assert inserted_songs[0]['album_json']
175 assert '"id": 20' in inserted_songs[0]['album_json']
176 assert '"title": "专辑"' in inserted_songs[0]['album_json']
......
1 from unittest.mock import MagicMock, call 1 from unittest.mock import MagicMock, call
2 from etl_to_crawler.writer import ( 2 from etl_to_crawler.writer import (
3 upsert_kugou_albums,
4 upsert_kugou_singers,
5 upsert_kugou_songs,
6 upsert_netease_songs,
3 upsert_qq_singers, 7 upsert_qq_singers,
4 upsert_qq_singer_songs, 8 upsert_qq_singer_songs,
5 upsert_yinyan_song_records, 9 upsert_yinyan_song_records,
...@@ -45,3 +49,101 @@ def test_upsert_yinyan_song_records_replaces_existing_song_relation(): ...@@ -45,3 +49,101 @@ def test_upsert_yinyan_song_records_replaces_existing_song_relation():
45 assert delete_rows == [(10,), (11,)] 49 assert delete_rows == [(10,), (11,)]
46 assert 'INSERT INTO yinyan_song_records' in insert_sql 50 assert 'INSERT INTO yinyan_song_records' in insert_sql
47 assert insert_rows == [(10, 100), (11, 101)] 51 assert insert_rows == [(10, 100), (11, 101)]
52
53
54 def test_upsert_netease_songs_writes_album_json_column():
55 cur = MagicMock()
56 upsert_netease_songs(cur, [{
57 'song_uuid': 'song-uuid',
58 'platform_song_id': 300,
59 'album_id': 20,
60 'album_json': '{"id": 20, "title": "专辑"}',
61 'cover': 'https://example.com/cover.jpg',
62 'title': '标题',
63 'name': '词曲名',
64 'duration': 180,
65 'lyric': '歌词',
66 'composer_name': '曲作者',
67 'lyricist_name': '词作者',
68 'url': 'https://example.com/audio.mp3',
69 'lyric_url': 'https://example.com/lyric.lrc',
70 'platform_index_url': 'https://music.163.com/#/song?id=300',
71 'published_at': '2020-01-01',
72 'singers_json': '[]',
73 }])
74
75 sql, rows = cur.executemany.call_args[0]
76 assert 'album, singers' in sql
77 assert '%s::json, %s::jsonb' in sql
78 assert rows[0][-2] == '{"id": 20, "title": "专辑"}'
79
80
81 def test_upsert_kugou_songs_writes_provider_and_source_data():
82 cur = MagicMock()
83 upsert_kugou_songs(cur, [{
84 'song_uuid': 'song-uuid',
85 'platform_song_id': 200,
86 'hash': 'hash',
87 'album_audio_id': 300,
88 'album_id': 20,
89 'cover': 'https://example.com/cover.jpg',
90 'title': '标题',
91 'name': '词曲名',
92 'duration': 180,
93 'lyric': '歌词',
94 'composer_name': '曲作者',
95 'lyricist_name': '词作者',
96 'url': 'https://example.com/audio.mp3',
97 'lyric_url': 'https://example.com/lyric.lrc',
98 'platform_index_url': 'https://www.kugou.com/song/#hash=hash',
99 'published_at': '2020-01-01',
100 'singers_json': '[]',
101 'provider_name': 'kugou',
102 'crawler_source_data': '{"id": 200}',
103 }])
104
105 sql, rows = cur.executemany.call_args[0]
106 assert 'provider_name, crawler_source_data' in sql
107 assert '%s::json' in sql
108 assert rows[0][-2:] == ('kugou', '{"id": 200}')
109
110
111 def test_upsert_kugou_singers_writes_provider_and_source_data():
112 cur = MagicMock()
113 upsert_kugou_singers(cur, [{
114 'id': 1,
115 'name': '歌手',
116 'avatar': 'https://example.com/avatar.jpg',
117 'sex': 'M',
118 'area': '华语',
119 'index': 'G',
120 'intro': '简介',
121 'home_url': 'https://example.com/singer',
122 'provider_name': 'kugou',
123 'crawler_source_data': '{"id": 1}',
124 }])
125
126 sql, rows = cur.executemany.call_args[0]
127 assert 'provider_name, crawler_source_data' in sql
128 assert rows[0][-2:] == ('kugou', '{"id": 1}')
129
130
131 def test_upsert_kugou_albums_writes_provider_and_source_data():
132 cur = MagicMock()
133 upsert_kugou_albums(cur, [{
134 'id': 20,
135 'cover': 'https://example.com/album.jpg',
136 'title': '专辑',
137 'intro': '简介',
138 'type': '专辑类型',
139 'company_id': 7,
140 'company': '唱片公司',
141 'is_owner': 1,
142 'published_at': '2020-01-01',
143 'provider_name': 'kugou',
144 'crawler_source_data': '{"id": 20}',
145 }])
146
147 sql, rows = cur.executemany.call_args[0]
148 assert 'provider_name, crawler_source_data' in sql
149 assert rows[0][-2:] == ('kugou', '{"id": 20}')
......