Commit 443980dc 443980dc1b929d6398b24b3e4159819a7cc502e2 by 沈秋雨

refactor(lyric): 优化歌词转存逻辑及更新冲突处理

1 parent 55c6e4fb
......@@ -100,28 +100,55 @@ def _safe_transfer_audio(url, oss_key, bucket, base_url) -> tuple[str, str]:
return '', '' # 不将未转存的音频外链写入 crawler
def _safe_upload_lyric(
def _safe_prepare_lyric(
platform: str,
unique_id: str,
spider_lyric: str | None,
hk_lyrics_url: str | None,
bucket,
base_url: str,
) -> str:
"""优先将音眼歌词 URL 转存到 crawler OSS;spider 歌词文本仅作回退。"""
) -> tuple[str, str]:
"""返回实际写入的 (lyric_url, lyric_text),音眼歌词优先、spider 回退。"""
candidates: list[tuple[str, str]] = []
if str(hk_lyrics_url or '').strip():
try:
hk_lyric = download_text_url(hk_lyrics_url, base_url)
uploaded_url = upload_plain_lyric_to_bucket(platform, unique_id, hk_lyric, bucket, base_url)
if str(hk_lyric or '').strip():
candidates.append(('HK', hk_lyric))
except Exception as e:
log.warning('HK lyric download failed for %s/%s: %s', platform, unique_id, e)
spider_text = spider_lyric or ''
if str(spider_text).strip() and all(text != spider_text for _, text in candidates):
candidates.append(('spider', spider_text))
for source, lyric_text in candidates:
try:
uploaded_url = upload_plain_lyric_to_bucket(
platform, unique_id, lyric_text, bucket, base_url,
)
if uploaded_url:
return uploaded_url
return uploaded_url, lyric_text
except Exception as e:
log.warning('HK lyric transfer failed for %s/%s: %s', platform, unique_id, e)
try:
return upload_plain_lyric_to_bucket(platform, unique_id, spider_lyric or '', bucket, base_url)
except Exception as e:
log.warning("Lyric upload failed for %s/%s: %s", platform, unique_id, e)
return ''
log.warning('%s lyric upload failed for %s/%s: %s', source, platform, unique_id, e)
# records2 允许 URL 转存失败时入库;若文本已成功获取,仍尽量保留文本。
return '', candidates[0][1] if candidates else ''
def _safe_upload_lyric(
platform: str,
unique_id: str,
spider_lyric: str | None,
hk_lyrics_url: str | None,
bucket,
base_url: str,
) -> str:
"""兼容仅需要歌词 URL 的调用方。"""
lyric_url, _ = _safe_prepare_lyric(
platform, unique_id, spider_lyric, hk_lyrics_url, bucket, base_url,
)
return lyric_url
def _required_import_assets() -> tuple[str, ...]:
......@@ -134,7 +161,9 @@ def _required_import_assets() -> tuple[str, ...]:
def _require_primary_assets(assets: dict, required: tuple[str, ...] | None = None) -> None:
"""校验当前导入支线要求的资源,允许非必需资源以空值入库。"""
audio_url, _ = assets.get('audio', ('', ''))
values = {'audio': audio_url, 'cover': assets.get('cover'), 'lyric': assets.get('lyric')}
lyric_asset = assets.get('lyric')
lyric_url = lyric_asset[0] if isinstance(lyric_asset, tuple) else lyric_asset
values = {'audio': audio_url, 'cover': assets.get('cover'), 'lyric': lyric_url}
required = required or ('audio', 'cover', 'lyric')
missing = [name for name in required if not values.get(name)]
if missing:
......@@ -330,7 +359,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
build_oss_key('qq', 'cover', str(song_id_int) + '.jpg'),
bucket, base_url,
),
'lyric': lambda: _safe_upload_lyric('qq', mid, raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('qq', mid, raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -348,7 +377,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
_require_primary_assets(assets, _required_import_assets())
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
singer_rows = [{
......@@ -422,7 +451,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......@@ -460,7 +489,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s
build_oss_key('kugou', 'cover', str(song_id) + '.jpg'),
bucket, base_url,
),
'lyric': lambda: _safe_upload_lyric('kugou', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('kugou', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -478,7 +507,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s
_require_primary_assets(assets, _required_import_assets())
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
singer_rows = [{
......@@ -535,7 +564,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......@@ -572,7 +601,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict,
build_oss_key('netease', 'cover', str(song_id) + '.jpg'),
bucket, base_url,
),
'lyric': lambda: _safe_upload_lyric('netease', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('netease', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -590,7 +619,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict,
_require_primary_assets(assets, _required_import_assets())
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
singer_rows = [{
......@@ -659,7 +688,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict,
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......@@ -870,7 +899,7 @@ def _process_qq(
build_oss_key('qq', 'cover', str(song_id_int) + '.jpg'),
bucket, base_url
),
'lyric': lambda: _safe_upload_lyric('qq', mid, raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('qq', mid, raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -890,7 +919,7 @@ def _process_qq(
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
# 歌手头像转移 + 写入 singers
......@@ -967,7 +996,7 @@ def _process_qq(
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......@@ -1026,7 +1055,7 @@ def _process_kugou(
build_oss_key('kugou', 'cover', str(song_id) + '.jpg'),
bucket, base_url
),
'lyric': lambda: _safe_upload_lyric('kugou', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('kugou', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -1046,7 +1075,7 @@ def _process_kugou(
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
singer_rows = []
......@@ -1103,7 +1132,7 @@ def _process_kugou(
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......@@ -1161,7 +1190,7 @@ def _process_netease(
build_oss_key('netease', 'cover', str(song_id) + '.jpg'),
bucket, base_url
),
'lyric': lambda: _safe_upload_lyric('netease', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
'lyric': lambda: _safe_prepare_lyric('netease', str(song_id), raw_lyric, hk_row.get('lyrics_url'), bucket, base_url),
}
if sp.get('album_id') and sp.get('album_cover'):
tasks['album_cover'] = lambda: _safe_transfer(
......@@ -1181,7 +1210,7 @@ def _process_netease(
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
lyric_url, lyric_text = assets['lyric']
album_cover = assets.get('album_cover') or cover_url
singer_rows = []
......@@ -1250,7 +1279,7 @@ def _process_netease(
'version': version,
'name': hk_row['name'],
'duration': sp.get('duration') or hk_row.get('song_time') or 0,
'lyric': ensure_newlines(raw_lyric),
'lyric': ensure_newlines(lyric_text),
'composer_name': sp.get('composer_name') or hk_row.get('composer'),
'lyricist_name': sp.get('lyricist_name') or hk_row.get('lyricist'),
'url': audio_url,
......
......@@ -442,7 +442,14 @@ def upsert_qq_songs(cur, songs: list[dict]) -> None:
platform_index_url, published_at, album, singers, status, created_at, updated_at,
version, audio_md5, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW(), %s, %s, %s, %s::json)
ON CONFLICT (platform_song_id) DO NOTHING
ON CONFLICT (platform_song_id) DO UPDATE SET
cover = CASE WHEN BTRIM(COALESCE(crawler_qqmusic_songs.cover, '')) = ''
THEN EXCLUDED.cover ELSE crawler_qqmusic_songs.cover END,
lyric = CASE WHEN BTRIM(COALESCE(crawler_qqmusic_songs.lyric, '')) = ''
THEN EXCLUDED.lyric ELSE crawler_qqmusic_songs.lyric END,
lyric_url = CASE WHEN BTRIM(COALESCE(crawler_qqmusic_songs.lyric_url, '')) = ''
THEN EXCLUDED.lyric_url ELSE crawler_qqmusic_songs.lyric_url END,
updated_at = NOW()
"""
rows = [(
s['song_uuid'], s['platform_song_id'], s['mid'], s.get('album_id'),
......@@ -534,7 +541,14 @@ def upsert_kugou_songs(cur, songs: list[dict]) -> None:
platform_index_url, published_at, singers, status, created_at, updated_at,
version, audio_md5, provider_name, crawler_source_data)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW(), %s, %s, %s, %s::json)
ON CONFLICT (platform_song_id) DO NOTHING
ON CONFLICT (platform_song_id) DO UPDATE SET
cover = CASE WHEN BTRIM(COALESCE(crawler_kugou_songs.cover, '')) = ''
THEN EXCLUDED.cover ELSE crawler_kugou_songs.cover END,
lyric = CASE WHEN BTRIM(COALESCE(crawler_kugou_songs.lyric, '')) = ''
THEN EXCLUDED.lyric ELSE crawler_kugou_songs.lyric END,
lyric_url = CASE WHEN BTRIM(COALESCE(crawler_kugou_songs.lyric_url, '')) = ''
THEN EXCLUDED.lyric_url ELSE crawler_kugou_songs.lyric_url END,
updated_at = NOW()
"""
rows = [(
s['song_uuid'], s['platform_song_id'], s.get('hash', ''),
......@@ -623,7 +637,14 @@ def upsert_netease_songs(cur, songs: list[dict]) -> None:
platform_index_url, published_at, album, singers, status, created_at, updated_at,
version, audio_md5, 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, %s, %s::json)
ON CONFLICT (platform_song_id) DO NOTHING
ON CONFLICT (platform_song_id) DO UPDATE SET
cover = CASE WHEN BTRIM(COALESCE(crawler_netease_songs.cover, '')) = ''
THEN EXCLUDED.cover ELSE crawler_netease_songs.cover END,
lyric = CASE WHEN BTRIM(COALESCE(crawler_netease_songs.lyric, '')) = ''
THEN EXCLUDED.lyric ELSE crawler_netease_songs.lyric END,
lyric_url = CASE WHEN BTRIM(COALESCE(crawler_netease_songs.lyric_url, '')) = ''
THEN EXCLUDED.lyric_url ELSE crawler_netease_songs.lyric_url END,
updated_at = NOW()
"""
rows = [(
s['song_uuid'], s['platform_song_id'], s.get('album_id'),
......
......@@ -105,32 +105,35 @@ def test_records2_import_allows_empty_cover_and_lyric(monkeypatch):
}, runner._required_import_assets())
def test_lyric_upload_prefers_hk_url_and_forces_oss_transfer(monkeypatch):
def test_lyric_prepare_prefers_hk_text_for_both_url_and_text_column(monkeypatch):
download = MagicMock(return_value='[00:01.00]HK lyric')
monkeypatch.setattr(runner, 'download_text_url', download)
upload = MagicMock(return_value='https://archive.example/crawler/lyric/qq/mid.txt')
monkeypatch.setattr(runner, 'upload_plain_lyric_to_bucket', upload)
bucket = object()
result = runner._safe_upload_lyric(
result = runner._safe_prepare_lyric(
'qq', 'mid', 'spider lyric', 'https://source.example/lyric.txt', bucket, 'https://archive.example',
)
assert result == 'https://archive.example/crawler/lyric/qq/mid.txt'
assert result == (
'https://archive.example/crawler/lyric/qq/mid.txt',
'[00:01.00]HK lyric',
)
download.assert_called_once_with('https://source.example/lyric.txt', 'https://archive.example')
upload.assert_called_once_with(
'qq', 'mid', '[00:01.00]HK lyric', bucket, 'https://archive.example',
)
def test_lyric_upload_falls_back_to_spider_text_when_hk_transfer_fails(monkeypatch):
def test_lyric_prepare_falls_back_to_spider_text_when_hk_download_fails(monkeypatch):
monkeypatch.setattr(runner, 'download_text_url', MagicMock(side_effect=TimeoutError('timeout')))
upload = MagicMock(return_value='https://archive.example/crawler/lyric/qq/mid.txt')
monkeypatch.setattr(runner, 'upload_plain_lyric_to_bucket', upload)
assert runner._safe_upload_lyric(
assert runner._safe_prepare_lyric(
'qq', 'mid', 'spider lyric', 'https://source.example/lyric.txt', object(), 'https://archive.example',
) == 'https://archive.example/crawler/lyric/qq/mid.txt'
) == ('https://archive.example/crawler/lyric/qq/mid.txt', 'spider lyric')
assert upload.call_args.args[2] == 'spider lyric'
......@@ -519,7 +522,10 @@ def test_process_qq_builds_album_json_for_song_insert(monkeypatch):
monkeypatch.setattr(runner, 'fetch_qq_singers', lambda conn, song_ids: {})
monkeypatch.setattr(runner, '_safe_transfer', lambda url, oss_key, bucket, base_url: url)
monkeypatch.setattr(runner, '_safe_transfer_audio', lambda url, oss_key, bucket, base_url: (url, 'audio-md5'))
monkeypatch.setattr(runner, '_safe_upload_lyric', lambda *args: 'https://bucket.example.com/lyric.txt')
monkeypatch.setattr(
runner, '_safe_prepare_lyric',
lambda *args: ('https://bucket.example.com/lyric.txt', '[00:01.00]HK歌词'),
)
monkeypatch.setattr(runner, 'upsert_qq_singers', lambda cur, singers: None)
monkeypatch.setattr(runner, 'upsert_qq_albums', lambda cur, albums: None)
monkeypatch.setattr(runner, 'upsert_qq_songs', lambda cur, songs: inserted_songs.extend(songs))
......@@ -550,6 +556,8 @@ def test_process_qq_builds_album_json_for_song_insert(monkeypatch):
assert '"title": "专辑"' in inserted_songs[0]['album_json']
assert inserted_songs[0]['title'] == '化风行万里'
assert inserted_songs[0]['version'] == 'DJ默涵版'
assert inserted_songs[0]['lyric'] == '[00:01.00]HK歌词'
assert inserted_songs[0]['lyric_url'] == 'https://bucket.example.com/lyric.txt'
def test_process_netease_builds_album_json_for_song_insert(monkeypatch):
......@@ -582,7 +590,10 @@ def test_process_netease_builds_album_json_for_song_insert(monkeypatch):
monkeypatch.setattr(runner, 'fetch_netease_singers', lambda conn, song_ids: {})
monkeypatch.setattr(runner, '_safe_transfer', lambda url, oss_key, bucket, base_url: url)
monkeypatch.setattr(runner, '_safe_transfer_audio', lambda url, oss_key, bucket, base_url: (url, 'audio-md5'))
monkeypatch.setattr(runner, '_safe_upload_lyric', lambda *args: 'https://bucket.example.com/lyric.txt')
monkeypatch.setattr(
runner, '_safe_prepare_lyric',
lambda *args: ('https://bucket.example.com/lyric.txt', '[00:01.00]HK歌词'),
)
monkeypatch.setattr(runner, 'upsert_netease_singers', lambda cur, singers: None)
monkeypatch.setattr(runner, 'upsert_netease_albums', lambda cur, albums: None)
monkeypatch.setattr(runner, 'upsert_netease_songs', lambda cur, songs: inserted_songs.extend(songs))
......@@ -625,7 +636,10 @@ def test_process_netease_uses_prefetched_song_and_singers(monkeypatch):
monkeypatch.setattr(runner, 'fetch_netease_singers', fetch_singers)
monkeypatch.setattr(runner, '_safe_transfer', lambda url, oss_key, bucket, base_url: url)
monkeypatch.setattr(runner, '_safe_transfer_audio', lambda url, oss_key, bucket, base_url: (url, 'audio-md5'))
monkeypatch.setattr(runner, '_safe_upload_lyric', lambda *args: 'https://bucket.example.com/lyric.txt')
monkeypatch.setattr(
runner, '_safe_prepare_lyric',
lambda *args: ('https://bucket.example.com/lyric.txt', '[00:01.00]HK歌词'),
)
monkeypatch.setattr(runner, 'upsert_netease_singers', lambda cur, singers: None)
monkeypatch.setattr(runner, 'upsert_netease_albums', lambda cur, albums: None)
monkeypatch.setattr(runner, 'upsert_netease_songs', lambda cur, songs: inserted_songs.extend(songs))
......
from unittest.mock import MagicMock, call
import pytest
import etl_to_crawler.writer as writer
from etl_to_crawler.writer import (
fetch_pending_yinyan_song_records,
......@@ -327,6 +329,34 @@ def test_upsert_kugou_songs_writes_provider_and_source_data():
assert rows[0][-4:] == ('DJ默涵版', 'md5-200', 'yinyan', '{"id": 200}')
@pytest.mark.parametrize(
('upsert', 'table_name', 'extra'),
[
(upsert_qq_songs, 'crawler_qqmusic_songs', {'mid': 'song-mid'}),
(upsert_kugou_songs, 'crawler_kugou_songs', {'hash': 'hash'}),
(upsert_netease_songs, 'crawler_netease_songs', {}),
],
)
def test_song_conflict_only_fills_empty_cover_and_lyric_fields(upsert, table_name, extra):
cur = MagicMock()
upsert(cur, [{
'song_uuid': 'song-uuid',
'platform_song_id': 200,
'name': '歌名',
'cover': 'https://example.com/cover.jpg',
'lyric': '歌词文本',
'lyric_url': 'https://example.com/lyric.txt',
**extra,
}])
sql = cur.executemany.call_args.args[0]
assert 'ON CONFLICT (platform_song_id) DO UPDATE SET' in sql
assert f"BTRIM(COALESCE({table_name}.cover, '')) = ''" in sql
assert f"BTRIM(COALESCE({table_name}.lyric, '')) = ''" in sql
assert f"BTRIM(COALESCE({table_name}.lyric_url, '')) = ''" in sql
assert 'THEN EXCLUDED.lyric' in sql
def test_upsert_kugou_singers_writes_provider_and_source_data():
cur = MagicMock()
upsert_kugou_singers(cur, [{
......