Commit 8bd8c02b 8bd8c02b2d435d6a7d7ae5b91dbed2876e9918c3 by 沈秋雨

refactor(crawler): 统一新增 provider_name 与 crawler_source_data 字段

- 替换所有歌手、专辑和歌曲的 provider_name 常量为 PROVIDER_YINYAN
- 增加 crawler_source_data 字段传递原始采集数据 JSON
- 更新 upsert 函数,支持写入 provider_name 和 crawler_source_data 字段
- 修改数据库写入 SQL,新增 provider_name 和 crawler_source_data 两列
- 添加相关单元测试,验证 provider_name 和 crawler_source_data 是否正确写入
- 统一整合 QQ、酷狗、网易等平台数据结构,增强数据来源追踪能力
1 parent ac91b0fe
......@@ -26,7 +26,7 @@ from .lyric import strip_timestamps
logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
log = logging.getLogger(__name__)
PROVIDER_KUGOU = 'kugou'
PROVIDER_YINYAN = 'yinyan'
def _safe_transfer(url, oss_key, bucket, base_url):
......@@ -83,7 +83,13 @@ def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
build_oss_key('qq', 'singer', sg['mid'] + '.jpg'),
bucket, base_url
)
singer_rows.append({**sg, 'id': sg['singer_id'], 'avatar': avatar})
singer_rows.append({
**sg,
'id': sg['singer_id'],
'avatar': avatar,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sg),
})
upsert_qq_singers(pg_cur, singer_rows)
# singers JSONB
......@@ -103,6 +109,19 @@ def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
'intro': sp.get('album_intro'), 'type': sp.get('album_type') or '',
'company_id': sp.get('company_id') or 0, 'company': sp.get('company') or '',
'is_owner': sp.get('is_owner') or 0, 'published_at': sp.get('album_published_at'),
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json({
'id': sp.get('album_id'),
'mid': sp.get('album_mid'),
'cover': sp.get('album_cover'),
'title': sp.get('album_title'),
'intro': sp.get('album_intro'),
'type': sp.get('album_type'),
'company_id': sp.get('company_id'),
'company': sp.get('company'),
'is_owner': sp.get('is_owner'),
'published_at': sp.get('album_published_at'),
}),
}])
# 写入 song
......@@ -125,6 +144,8 @@ def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
'platform_index_url': sp.get('platform_index_url') or f'https://y.qq.com/n/ryqq/songDetail/{mid}',
'published_at': sp.get('published_at') or hk_row.get('issue_time'),
'singers_json': singers_json,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sp),
}])
# 查询实际 UUID(ON CONFLICT DO NOTHING 时使用已有 UUID)
pg_cur.execute('SELECT id FROM crawler_qqmusic_songs WHERE platform_song_id = %s', (platform_song_id,))
......@@ -178,7 +199,7 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
**sg,
'id': sg['singer_id'],
'avatar': avatar,
'provider_name': PROVIDER_KUGOU,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sg),
})
upsert_kugou_singers(pg_cur, singer_rows)
......@@ -198,7 +219,7 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
'type': sp.get('album_type') or '', 'company_id': sp.get('company_id') or 0,
'company': sp.get('company') or '', 'is_owner': sp.get('is_owner') or 0,
'published_at': sp.get('album_published_at'),
'provider_name': PROVIDER_KUGOU,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json({
'id': sp.get('album_id'),
'cover': sp.get('album_cover'),
......@@ -231,7 +252,7 @@ def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url
'platform_index_url': sp.get('platform_index_url') or f'http://www.kugou.com/song/#hash={sp.get("hid") or pr.get("platform_mid", "")}',
'published_at': sp.get('published_at') or hk_row.get('issue_time'),
'singers_json': singers_json,
'provider_name': PROVIDER_KUGOU,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sp),
}])
# 查询实际 UUID(ON CONFLICT DO NOTHING 时使用已有 UUID)
......@@ -281,7 +302,13 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u
build_oss_key('netease', 'singer', str(sg['singer_id']) + '.jpg'),
bucket, base_url
)
singer_rows.append({**sg, 'id': sg['singer_id'], 'avatar': avatar})
singer_rows.append({
**sg,
'id': sg['singer_id'],
'avatar': avatar,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sg),
})
upsert_netease_singers(pg_cur, singer_rows)
singers_json = json.dumps([{
......@@ -312,6 +339,18 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u
'type': sp.get('album_type') or '', 'company_id': sp.get('company_id') or 0,
'company': sp.get('company') or '', 'is_owner': sp.get('is_owner') or 0,
'published_at': sp.get('album_published_at'),
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json({
'id': sp.get('album_id'),
'cover': sp.get('album_cover'),
'title': sp.get('album_title'),
'intro': sp.get('album_intro'),
'type': sp.get('album_type'),
'company_id': sp.get('company_id'),
'company': sp.get('company'),
'is_owner': sp.get('is_owner'),
'published_at': sp.get('album_published_at'),
}),
}])
song_uuid = str(uuid.uuid4())
......@@ -332,6 +371,8 @@ def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_u
'platform_index_url': sp.get('platform_index_url') or f'https://music.163.com/#/song?id={song_id}',
'published_at': sp.get('published_at') or hk_row.get('issue_time'),
'singers_json': singers_json,
'provider_name': PROVIDER_YINYAN,
'crawler_source_data': _source_json(sp),
}])
# 查询实际 UUID(ON CONFLICT DO NOTHING 时使用已有 UUID)
pg_cur.execute('SELECT id FROM crawler_netease_songs WHERE platform_song_id = %s', (song_id,))
......
......@@ -25,14 +25,17 @@ def upsert_qq_singers(cur, singers: list[dict]) -> None:
return
sql = """
INSERT INTO crawler_qqmusic_singers
(id, mid, name, avatar, sex, area, "index", intro, home_url, created_at, updated_at)
VALUES (%s, %s, %s, %s, %s::singer_sex, %s::singer_area, %s::singer_index, %s, %s, NOW(), NOW())
(id, mid, name, avatar, sex, area, "index", intro, home_url,
created_at, updated_at, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s::singer_sex, %s::singer_area, %s::singer_index, %s, %s,
NOW(), NOW(), %s, %s::json)
ON CONFLICT (id) DO NOTHING
"""
rows = [(
s['id'], s['mid'], s['name'], s.get('avatar', ''),
s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
s.get('intro'), s.get('home_url'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in singers]
cur.executemany(sql, rows)
......@@ -42,14 +45,16 @@ def upsert_qq_albums(cur, albums: list[dict]) -> None:
return
sql = """
INSERT INTO crawler_qqmusic_albums
(id, mid, cover, title, intro, type, company_id, company, is_owner, published_at, created_at, updated_at)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW())
(id, mid, cover, title, intro, type, company_id, company, is_owner, published_at,
created_at, updated_at, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
ON CONFLICT (id) DO NOTHING
"""
rows = [(
a['id'], a.get('mid', ''), a.get('cover', ''), a.get('title', ''),
a.get('intro'), a.get('type', ''), a.get('company_id', 0) or 0,
a.get('company', ''), a.get('is_owner', 0) or 0, a.get('published_at'),
a.get('provider_name'), a.get('crawler_source_data'),
) for a in albums]
cur.executemany(sql, rows)
......@@ -61,8 +66,9 @@ def upsert_qq_songs(cur, songs: list[dict]) -> None:
INSERT INTO crawler_qqmusic_songs
(id, platform_song_id, mid, album_id, cover, title, name, duration,
lyric, composer_name, lyricist_name, url, lyric_url,
platform_index_url, published_at, singers, status, created_at, updated_at)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW())
platform_index_url, published_at, singers, status, created_at, updated_at,
provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW(), %s, %s::json)
ON CONFLICT (platform_song_id) DO NOTHING
"""
rows = [(
......@@ -71,6 +77,7 @@ def upsert_qq_songs(cur, songs: list[dict]) -> None:
s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'),
s.get('url', ''), s.get('lyric_url'),
s.get('platform_index_url'), s.get('published_at'), s.get('singers_json', '[]'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in songs]
cur.executemany(sql, rows)
......@@ -195,14 +202,17 @@ def upsert_netease_singers(cur, singers: list[dict]) -> None:
return
sql = """
INSERT INTO crawler_netease_singers
(id, name, avatar, sex, area, "index", intro, home_url, created_at, updated_at)
VALUES (%s, %s, %s, %s::netease_singer_sex, %s::netease_singer_area, %s::netease_singer_index, %s, %s, NOW(), NOW())
(id, name, avatar, sex, area, "index", intro, home_url,
created_at, updated_at, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s::netease_singer_sex, %s::netease_singer_area, %s::netease_singer_index, %s, %s,
NOW(), NOW(), %s, %s::json)
ON CONFLICT (id) DO NOTHING
"""
rows = [(
s['id'], s['name'], s.get('avatar', ''),
s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
s.get('intro'), s.get('home_url'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in singers]
cur.executemany(sql, rows)
......@@ -212,14 +222,16 @@ def upsert_netease_albums(cur, albums: list[dict]) -> None:
return
sql = """
INSERT INTO crawler_netease_albums
(id, cover, title, intro, type, company_id, company, is_owner, published_at, created_at, updated_at)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW())
(id, cover, title, intro, type, company_id, company, is_owner, published_at,
created_at, updated_at, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
ON CONFLICT (id) DO NOTHING
"""
rows = [(
a['id'], a.get('cover', ''), a.get('title', ''), a.get('intro'),
a.get('type', ''), a.get('company_id', 0) or 0,
a.get('company', ''), a.get('is_owner', 0) or 0, a.get('published_at'),
a.get('provider_name'), a.get('crawler_source_data'),
) for a in albums]
cur.executemany(sql, rows)
......@@ -231,8 +243,9 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None:
INSERT INTO crawler_netease_songs
(id, platform_song_id, album_id, cover, title, name, duration,
lyric, composer_name, lyricist_name, url, lyric_url,
platform_index_url, published_at, album, singers, status, created_at, updated_at)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW())
platform_index_url, published_at, album, singers, status, created_at, updated_at,
provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW(), %s, %s::json)
ON CONFLICT (platform_song_id) DO NOTHING
"""
rows = [(
......@@ -241,6 +254,7 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None:
s.get('lyric'), s.get('composer_name'), s.get('lyricist_name'),
s.get('url', ''), s.get('lyric_url'),
s.get('platform_index_url'), s.get('published_at'), s.get('album_json'), s.get('singers_json', '[]'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in songs]
cur.executemany(sql, rows)
......
......@@ -3,8 +3,12 @@ from etl_to_crawler.writer import (
upsert_kugou_albums,
upsert_kugou_singers,
upsert_kugou_songs,
upsert_netease_albums,
upsert_netease_singers,
upsert_netease_songs,
upsert_qq_albums,
upsert_qq_singers,
upsert_qq_songs,
upsert_qq_singer_songs,
upsert_yinyan_song_records,
)
......@@ -75,7 +79,7 @@ def test_upsert_netease_songs_writes_album_json_column():
sql, rows = cur.executemany.call_args[0]
assert 'album, singers' in sql
assert '%s::json, %s::jsonb' in sql
assert rows[0][-2] == '{"id": 20, "title": "专辑"}'
assert rows[0][14] == '{"id": 20, "title": "专辑"}'
def test_upsert_kugou_songs_writes_provider_and_source_data():
......@@ -98,14 +102,14 @@ def test_upsert_kugou_songs_writes_provider_and_source_data():
'platform_index_url': 'https://www.kugou.com/song/#hash=hash',
'published_at': '2020-01-01',
'singers_json': '[]',
'provider_name': 'kugou',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 200}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert '%s::json' in sql
assert rows[0][-2:] == ('kugou', '{"id": 200}')
assert rows[0][-2:] == ('yinyan', '{"id": 200}')
def test_upsert_kugou_singers_writes_provider_and_source_data():
......@@ -119,13 +123,13 @@ def test_upsert_kugou_singers_writes_provider_and_source_data():
'index': 'G',
'intro': '简介',
'home_url': 'https://example.com/singer',
'provider_name': 'kugou',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 1}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('kugou', '{"id": 1}')
assert rows[0][-2:] == ('yinyan', '{"id": 1}')
def test_upsert_kugou_albums_writes_provider_and_source_data():
......@@ -140,10 +144,136 @@ def test_upsert_kugou_albums_writes_provider_and_source_data():
'company': '唱片公司',
'is_owner': 1,
'published_at': '2020-01-01',
'provider_name': 'kugou',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 20}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('kugou', '{"id": 20}')
assert rows[0][-2:] == ('yinyan', '{"id": 20}')
def test_upsert_qq_entities_write_provider_and_source_data():
cur = MagicMock()
upsert_qq_singers(cur, [{
'id': 1,
'mid': 'singer-mid',
'name': '歌手',
'avatar': 'https://example.com/avatar.jpg',
'sex': 'M',
'area': '华语',
'index': 'G',
'intro': '简介',
'home_url': 'https://example.com/singer',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 1}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 1}')
cur = MagicMock()
upsert_qq_albums(cur, [{
'id': 20,
'mid': 'album-mid',
'cover': 'https://example.com/album.jpg',
'title': '专辑',
'intro': '简介',
'type': '专辑类型',
'company_id': 7,
'company': '唱片公司',
'is_owner': 1,
'published_at': '2020-01-01',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 20}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 20}')
cur = MagicMock()
upsert_qq_songs(cur, [{
'song_uuid': 'song-uuid',
'platform_song_id': 200,
'mid': 'song-mid',
'album_id': 20,
'cover': 'https://example.com/cover.jpg',
'title': '标题',
'name': '词曲名',
'duration': 180,
'lyric': '歌词',
'composer_name': '曲作者',
'lyricist_name': '词作者',
'url': 'https://example.com/audio.mp3',
'lyric_url': 'https://example.com/lyric.lrc',
'platform_index_url': 'https://y.qq.com/n/ryqq/songDetail/song-mid',
'published_at': '2020-01-01',
'singers_json': '[]',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 200}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 200}')
def test_upsert_netease_entities_write_provider_and_source_data():
cur = MagicMock()
upsert_netease_singers(cur, [{
'id': 1,
'name': '歌手',
'avatar': 'https://example.com/avatar.jpg',
'sex': 'M',
'area': '华语',
'index': 'G',
'intro': '简介',
'home_url': 'https://example.com/singer',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 1}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 1}')
cur = MagicMock()
upsert_netease_albums(cur, [{
'id': 20,
'cover': 'https://example.com/album.jpg',
'title': '专辑',
'intro': '简介',
'type': '专辑类型',
'company_id': 7,
'company': '唱片公司',
'is_owner': 1,
'published_at': '2020-01-01',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 20}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 20}')
cur = MagicMock()
upsert_netease_songs(cur, [{
'song_uuid': 'song-uuid',
'platform_song_id': 300,
'album_id': 20,
'album_json': '{"id": 20, "title": "专辑"}',
'cover': 'https://example.com/cover.jpg',
'title': '标题',
'name': '词曲名',
'duration': 180,
'lyric': '歌词',
'composer_name': '曲作者',
'lyricist_name': '词作者',
'url': 'https://example.com/audio.mp3',
'lyric_url': 'https://example.com/lyric.lrc',
'platform_index_url': 'https://music.163.com/#/song?id=300',
'published_at': '2020-01-01',
'singers_json': '[]',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 300}',
}])
sql, rows = cur.executemany.call_args[0]
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-2:] == ('yinyan', '{"id": 300}')
......