Commit 2911f147 2911f1479ec790b0d5aac60dfde250f4bcbb8313 by 沈秋雨

feat(crawler): 增加QQ无效专辑封面回填功能

- 在 run_etl.py 添加 backfill-qq-invalid-covers 命令行选项
- 配置文件中新增 HTTP_TRANSFER_RETRIES 和 RETRY_BACKOFF_SECONDS
- oss.py 中实现下载资源的重试逻辑,避免超时和连接断开失败
- runner.py 中引入资源必要性校验,确保音频、封面、歌词均成功转存
- 实现根据歌词和封面有效性选择更合适的QQ录音记录进行替换
- 实现 backfill_qq_invalid_covers 函数,处理无专辑占位封面的QQ歌曲
- 无候选的QQ歌曲在crawler和yinyan记录中删除,并软删除 HK 歌曲
- 资源转存失败的记录在yinyan_song_records标记,避免重复处理
- run函数中增加对QQ、酷狗、网易备选录音的支持和筛选
- 增加对应的数据库操作函数支持替换、删除以及状态标记
- 测试覆盖下载重试、资源转存失败处理及选择替代录音逻辑
- 所有下载失败均返回空URL,不阻断流程以保证稳定性
- 异常处理增强,确保未成功上传资源的歌曲不入库
- 未匹配到HK歌曲的情况增加日志警告及延迟重试机制
1 parent a1e6ab0d
......@@ -51,4 +51,6 @@ PLATFORMS = [PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE]
BATCH_SIZE = 100
BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000'))
HTTP_POOL_MAXSIZE = int(os.environ.get('HTTP_POOL_MAXSIZE', '128'))
HTTP_TRANSFER_RETRIES = int(os.environ.get('HTTP_TRANSFER_RETRIES', '3'))
HTTP_TRANSFER_RETRY_BACKOFF_SECONDS = float(os.environ.get('HTTP_TRANSFER_RETRY_BACKOFF_SECONDS', '1'))
OSS_CONNECTION_POOL_SIZE = int(os.environ.get('OSS_CONNECTION_POOL_SIZE', str(HTTP_POOL_MAXSIZE)))
......
import requests
import oss2
import math
import time
from urllib.parse import urlparse
from requests.adapters import HTTPAdapter
from .config import HTTP_POOL_MAXSIZE, OSS_CONFIG
from .config import (
HTTP_POOL_MAXSIZE, HTTP_TRANSFER_RETRIES, HTTP_TRANSFER_RETRY_BACKOFF_SECONDS, OSS_CONFIG,
)
from .utils import compute_audio_md5
_HTTP_SESSION = requests.Session()
......@@ -47,6 +50,25 @@ def _download_url(url: str, base_url: str) -> str:
return f"{download_base_url.rstrip('/')}{url[len(rewrite_from.rstrip('/')):]}"
def _download_content(url: str, base_url: str) -> bytes:
"""下载资源;针对超时、连接中断和服务端错误进行有限重试。"""
download_url = _download_url(url, base_url)
attempts = max(1, HTTP_TRANSFER_RETRIES)
for attempt in range(attempts):
try:
resp = _http_get(download_url, timeout=30)
resp.raise_for_status()
return resp.content
except requests.RequestException as exc:
status_code = getattr(getattr(exc, 'response', None), 'status_code', None)
# 4xx(429 除外)是确定性失败,不浪费时间重试。
if status_code is not None and 400 <= status_code < 500 and status_code != 429:
raise
if attempt == attempts - 1:
raise
time.sleep(HTTP_TRANSFER_RETRY_BACKOFF_SECONDS * (2 ** attempt))
def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> str:
"""
将 url 指向的文件转移到 archive-dev OSS 的 oss_key 路径。
......@@ -58,9 +80,8 @@ def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: s
return ''
if _is_target_oss_url(url, base_url):
return url
resp = _http_get(_download_url(url, base_url), timeout=30)
resp.raise_for_status()
bucket.put_object(oss_key, resp.content)
content = _download_content(url, base_url)
bucket.put_object(oss_key, content)
return f"{base_url.rstrip('/')}/{oss_key}"
......@@ -74,10 +95,9 @@ def transfer_url_with_md5(url: str | None, oss_key: str, bucket: oss2.Bucket, ba
return '', ''
if _is_target_oss_url(url, base_url):
return url, ''
resp = _http_get(_download_url(url, base_url), timeout=30)
resp.raise_for_status()
audio_md5 = compute_audio_md5(resp.content)
bucket.put_object(oss_key, resp.content)
content = _download_content(url, base_url)
audio_md5 = compute_audio_md5(content)
bucket.put_object(oss_key, content)
return f"{base_url.rstrip('/')}/{oss_key}", audio_md5
......
......@@ -111,6 +111,21 @@ def fetch_hk_songs_by_source_ids(conn: pymysql.Connection, source_song_ids: list
}
def mark_hk_songs_deleted(conn: pymysql.Connection, source_song_ids: list[int]) -> int:
"""将无法导入的源歌曲在 hk_songs_test 中软删除。"""
if not source_song_ids:
return 0
unique_ids = list(dict.fromkeys(int(song_id) for song_id in source_song_ids))
placeholders = ','.join(['%s'] * len(unique_ids))
with conn.cursor() as cur:
cur.execute(
f"UPDATE hk_songs_test SET deleted = '1' "
f"WHERE source_song_id IN ({placeholders}) AND deleted = '0'",
unique_ids,
)
return cur.rowcount
def fetch_platform_records(source_conn: pymysql.Connection, song_ids: list[int]) -> list[dict]:
if not song_ids:
return []
......
......@@ -9,6 +9,7 @@ from .connections import get_hk_songs_conn, get_source_conn, get_spider_conn, ge
from .reader import (
iter_hk_songs_batches,
fetch_hk_songs_by_source_ids,
mark_hk_songs_deleted,
fetch_platform_records,
fetch_all_platform_records,
fetch_platform_records_by_record_ids,
......@@ -36,6 +37,8 @@ from .writer import (
upsert_netease_singer_songs, upsert_netease_singer_albums,
fetch_qq_songs_missing_singers, fetch_kugou_songs_missing_singers, fetch_netease_songs_missing_singers,
update_qq_song_singers, update_kugou_song_singers, update_netease_song_singers,
fetch_qq_songs_with_invalid_covers, replace_qq_song, replace_yinyan_song_record,
delete_qq_song_and_yinyan_record, mark_yinyan_record_resource_failed,
)
from .oss import transfer_url, transfer_url_with_md5, build_oss_key
from .utils import split_title_version, upload_plain_lyric_to_bucket
......@@ -46,6 +49,11 @@ log = logging.getLogger(__name__)
PROVIDER_YINYAN = 'yinyan'
MAX_RESOURCE_WORKERS = 8
MAX_IMPORT_WORKERS = 16
QQ_MISSING_ALBUM_COVER = 'http://y.gtimg.cn/music/photo_new/T002R300x300M000.jpg'
class RequiredAssetTransferError(RuntimeError):
"""音频、封面或歌词未成功落到目标 OSS,当前歌曲不得入库。"""
def _run_io_tasks(tasks: dict) -> dict:
......@@ -65,7 +73,7 @@ def _safe_transfer(url, oss_key, bucket, base_url):
return transfer_url(url, oss_key, bucket, base_url)
except Exception as e:
log.warning("OSS transfer failed for %s: %s", url, e)
return url # 失败时保留原 URL,不阻断流程
return '' # 不将未转存的外链写入 crawler
def _safe_transfer_audio(url, oss_key, bucket, base_url) -> tuple[str, str]:
......@@ -73,15 +81,31 @@ def _safe_transfer_audio(url, oss_key, bucket, base_url) -> tuple[str, str]:
return transfer_url_with_md5(url, oss_key, bucket, base_url)
except Exception as e:
log.warning("Audio transfer failed for %s: %s", url, e)
return url or '', ''
return '', '' # 不将未转存的音频外链写入 crawler
def _safe_upload_lyric(platform: str, unique_id: str, lyric: str | None, fallback_url: str | None, bucket, base_url: str) -> str:
try:
return upload_plain_lyric_to_bucket(platform, unique_id, lyric or '', bucket, base_url) or (fallback_url or '')
uploaded_url = upload_plain_lyric_to_bucket(platform, unique_id, lyric or '', bucket, base_url)
if uploaded_url:
return uploaded_url
fallback_url = str(fallback_url or '')
return fallback_url if fallback_url.startswith(base_url.rstrip('/') + '/') else ''
except Exception as e:
log.warning("Lyric upload failed for %s/%s: %s", platform, unique_id, e)
return fallback_url or ''
return ''
def _require_primary_assets(assets: dict) -> None:
"""音频、歌曲封面、歌词均为入库必需资源,任一失败即拒绝整条记录。"""
audio_url, _ = assets.get('audio', ('', ''))
missing = [
name for name, value in (
('audio', audio_url), ('cover', assets.get('cover')), ('lyric', assets.get('lyric')),
) if not value
]
if missing:
raise RequiredAssetTransferError(f"required asset transfer failed: {', '.join(missing)}")
def _json_default(value):
......@@ -92,6 +116,123 @@ def _source_json(data: dict) -> str:
return json.dumps(data, ensure_ascii=False, default=_json_default)
def _is_usable_qq_cover(url: str | None) -> bool:
return bool(url and str(url).strip() and str(url).strip() != QQ_MISSING_ALBUM_COVER)
def _qq_cover_source(hk_row: dict, song_data: dict) -> str:
"""HK 封面优先;QQ 的 M000 占位图视为无封面,改用当前录音的封面。"""
hk_cover = hk_row.get('cover_url')
return hk_cover if _is_usable_qq_cover(hk_cover) else (song_data.get('cover') or '')
def _has_usable_qq_lyric(song_data: dict) -> bool:
return bool(str(song_data.get('lyric') or '').strip())
def _has_kugou_cover(hk_row: dict, song_data: dict) -> bool:
"""酷狗封面来源:hk_row.cover_url 优先,回退到 spider 录音封面。"""
return bool((hk_row.get('cover_url') or '').strip() or (song_data.get('cover') or '').strip())
def _has_kugou_lyric(song_data: dict) -> bool:
return bool(str(song_data.get('lyric') or '').strip())
def _has_netease_cover(hk_row: dict, song_data: dict) -> bool:
"""网易封面来源:hk_row.cover_url 优先,回退到 spider 录音封面。"""
return bool((hk_row.get('cover_url') or '').strip() or (song_data.get('cover') or '').strip())
def _has_netease_lyric(song_data: dict) -> bool:
return bool(str(song_data.get('lyric') or '').strip())
def _select_qq_import_record(
hk_row: dict,
current_record: dict,
candidate_records: list[dict],
songs_by_mid: dict[str, dict],
) -> dict | None:
"""为导入选择 QQ 录音;当前录音的歌词或封面不可用时,改选同源的合格候选。"""
current_song = songs_by_mid.get(current_record['platform_unique_key'])
if not current_song:
return None
needs_lyric = not _has_usable_qq_lyric(current_song)
needs_cover = not _is_usable_qq_cover(_qq_cover_source(hk_row, current_song)) or not current_song.get('album_id')
if not needs_lyric and not needs_cover:
return current_record
for candidate in candidate_records:
song = songs_by_mid.get(candidate['platform_unique_key'])
if not song:
continue
if needs_lyric and not _has_usable_qq_lyric(song):
continue
if needs_cover and (not song.get('album_id') or not _is_usable_qq_cover(song.get('cover'))):
continue
return candidate
return None
def _select_kugou_import_record(
hk_row: dict,
current_record: dict,
candidate_records: list[dict],
songs_by_id: dict[int, dict],
) -> dict | None:
"""为导入选择酷狗录音;当前录音的歌词或封面不可用时,改选同源的合格候选。"""
song_id = int(current_record['platform_unique_key'])
current_song = songs_by_id.get(song_id)
if not current_song:
return None
needs_lyric = not _has_kugou_lyric(current_song)
needs_cover = not _has_kugou_cover(hk_row, current_song)
if not needs_lyric and not needs_cover:
return current_record
for candidate in candidate_records:
cid = int(candidate['platform_unique_key'])
song = songs_by_id.get(cid)
if not song:
continue
if needs_lyric and not _has_kugou_lyric(song):
continue
if needs_cover and not _has_kugou_cover(hk_row, song):
continue
return candidate
return None
def _select_netease_import_record(
hk_row: dict,
current_record: dict,
candidate_records: list[dict],
songs_by_id: dict[int, dict],
) -> dict | None:
"""为导入选择网易录音;当前录音的歌词或封面不可用时,改选同源的合格候选。"""
song_id = int(current_record['platform_unique_key'])
current_song = songs_by_id.get(song_id)
if not current_song:
return None
needs_lyric = not _has_netease_lyric(current_song)
needs_cover = not _has_netease_cover(hk_row, current_song)
if not needs_lyric and not needs_cover:
return current_record
for candidate in candidate_records:
cid = int(candidate['platform_unique_key'])
song = songs_by_id.get(cid)
if not song:
continue
if needs_lyric and not _has_netease_lyric(song):
continue
if needs_cover and not _has_netease_cover(hk_row, song):
continue
return candidate
return None
def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, singer_list: list[dict]) -> dict:
mid = pr['platform_unique_key']
song_id_int = sp['id']
......@@ -101,7 +242,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
hk_row['audio_url'], build_oss_key('qq', 'audio', mid + '.mp3'), bucket, base_url
),
'cover': lambda: _safe_transfer(
hk_row.get('cover_url') or sp.get('cover', ''),
_qq_cover_source(hk_row, sp),
build_oss_key('qq', 'cover', str(song_id_int) + '.jpg'),
bucket, base_url,
),
......@@ -120,6 +261,7 @@ def _prepare_qq_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, sing
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -248,6 +390,7 @@ def _prepare_kugou_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict, s
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -359,6 +502,7 @@ def _prepare_netease_payload(hk_row: dict, pr: dict, bucket, base_url, sp: dict,
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -571,7 +715,7 @@ def _process_qq(
bucket, base_url
),
'cover': lambda: _safe_transfer(
hk_row.get('cover_url') or sp.get('cover', ''),
_qq_cover_source(hk_row, sp),
build_oss_key('qq', 'cover', str(song_id_int) + '.jpg'),
bucket, base_url
),
......@@ -592,6 +736,7 @@ def _process_qq(
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -747,6 +892,7 @@ def _process_kugou(
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -881,6 +1027,7 @@ def _process_netease(
)
)
assets = _run_io_tasks(tasks)
_require_primary_assets(assets)
audio_url, audio_md5 = assets['audio']
cover_url = assets['cover']
lyric_url = assets['lyric']
......@@ -1051,6 +1198,19 @@ def initialize_yinyan_song_records(platforms: list[str], max_batches: int | None
insert_yinyan_song_records(pg_cur, init_rows)
pg_conn.commit()
total += len(init_rows)
# 本批次中无法匹配到任何录音的歌曲 → 软删 hk_songs_test
matched_src_ids = {int(hk_row['source_song_id']) for hk_row in batch
if hk_row.get('source_song_id') and int(hk_row['source_song_id']) in pr_by_song}
unmatched_ids = [
int(hk_row['source_song_id']) for hk_row in batch
if hk_row.get('source_song_id') and int(hk_row['source_song_id']) not in matched_src_ids
and int(hk_row['source_song_id']) not in existing_song_ids
]
if unmatched_ids:
mark_hk_songs_deleted(hk_conn, unmatched_ids)
hk_conn.commit()
log.info("Init: marked %d unmatched songs as deleted", len(unmatched_ids))
finally:
hk_conn.close()
src_conn.close()
......@@ -1340,6 +1500,114 @@ def backfill_empty_singers(platforms: list[str], max_batches: int | None = None)
pg_conn.close()
def _is_target_oss_url(url: str | None, base_url: str) -> bool:
return bool(url and str(url).startswith(base_url.rstrip('/') + '/'))
def backfill_qq_invalid_covers(max_batches: int | None = None) -> None:
"""替换 QQ 无专辑占位封面的录音;无候选时删除 crawler 记录并软删 HK 歌曲。"""
hk_conn = get_hk_songs_conn()
src_conn = get_source_conn()
spider_conn = get_spider_conn()
pg_conn = get_pg_conn()
bucket = get_oss_bucket()
base_url = OSS_CONFIG['base_url']
batch_index = replaced = filtered = retry_later = 0
pbar = tqdm(desc='backfill-qq-invalid-covers')
try:
while max_batches is None or batch_index < max_batches:
with pg_conn.cursor() as pg_cur:
invalid_rows = fetch_qq_songs_with_invalid_covers(
pg_cur, QQ_MISSING_ALBUM_COVER, BATCH_SIZE,
)
if not invalid_rows:
break
source_ids = [row['song_id'] for row in invalid_rows]
hk_by_song = fetch_hk_songs_by_source_ids(hk_conn, source_ids)
records_by_song: dict[int, list[dict]] = {}
for record in fetch_all_platform_records(src_conn, source_ids):
if record['platform'] == PLATFORM_QQ:
records_by_song.setdefault(int(record['source_song_id']), []).append(record)
mids = list({record['platform_unique_key'] for records in records_by_song.values() for record in records})
songs_by_mid = fetch_qq_songs(spider_conn, mids) if mids else {}
qq_song_ids = [song['id'] for song in songs_by_mid.values()]
singers_by_song_id = fetch_qq_singers(spider_conn, qq_song_ids) if qq_song_ids else {}
replacements: list[tuple[dict, dict, dict]] = []
removals: list[dict] = []
for row in invalid_rows:
hk_row = hk_by_song.get(row['song_id'])
candidates = [
record for record in records_by_song.get(row['song_id'], [])
if (song := songs_by_mid.get(record['platform_unique_key']))
and song.get('album_id')
and _is_usable_qq_cover(song.get('cover'))
]
if not candidates:
removals.append(row)
continue
if not hk_row:
log.warning('QQ cover backfill cannot find hk_songs_test row: source_song_id=%s', row['song_id'])
retry_later += 1
continue
payload = None
selected_candidate = None
for candidate in candidates:
song = songs_by_mid[candidate['platform_unique_key']]
candidate_payload = _prepare_qq_payload(
hk_row, candidate, bucket, base_url, song,
singers_by_song_id.get(song['id'], []),
)
cover = candidate_payload['songs'][0]['cover']
if _is_target_oss_url(cover, base_url):
payload = candidate_payload
selected_candidate = candidate
break
log.warning(
'QQ cover replacement download failed: source_song_id=%s candidate_mid=%s',
row['song_id'], candidate['platform_unique_key'],
)
if payload is None:
retry_later += 1
continue
replacements.append((row, selected_candidate, payload))
with pg_conn.cursor() as pg_cur:
for row, candidate, payload in replacements:
song = payload['songs'][0]
upsert_qq_singers(pg_cur, payload['singers'])
upsert_qq_albums(pg_cur, payload['albums'])
replace_qq_song(pg_cur, row['platform_song_id'], song)
pg_cur.execute('DELETE FROM crawler_qqmusic_singer_songs WHERE song_id = %s', (row['id'],))
upsert_qq_singer_songs(pg_cur, [(sg['singer_id'], row['id']) for sg in payload['singers']])
upsert_qq_singer_albums(pg_cur, payload['singer_albums'])
replace_yinyan_song_record(
pg_cur, row['song_id'], row['platform_song_id'],
int(candidate['record_id']), int(song['platform_song_id']),
)
for row in removals:
delete_qq_song_and_yinyan_record(pg_cur, row['song_id'], row['platform_song_id'], row['id'])
pg_conn.commit()
if removals:
mark_hk_songs_deleted(hk_conn, [row['song_id'] for row in removals])
hk_conn.commit()
replaced += len(replacements)
filtered += len(removals)
batch_index += 1
pbar.update(1)
finally:
pbar.close()
hk_conn.close()
src_conn.close()
spider_conn.close()
pg_conn.close()
log.info('QQ invalid cover backfill: replaced=%d filtered=%d retry_later=%d', replaced, filtered, retry_later)
def run(
platforms: list[str],
max_batches: int | None = None,
......@@ -1372,10 +1640,37 @@ def run(
key = (int(pr['source_song_id']), int(pr['record_id']), pr['platform'])
pr_by_state_key[key] = pr
# 各平台在当前录音歌词或封面不可用时,需要同源的其他录音作为候选。
all_source_ids = [int(pending['song_id']) for pending in pending_records]
all_candidates = fetch_all_platform_records(src_conn, all_source_ids) if all_source_ids else []
qq_records_by_song: dict[int, list[dict]] = {}
kugou_records_by_song: dict[int, list[dict]] = {}
netease_records_by_song: dict[int, list[dict]] = {}
for candidate in all_candidates:
sid = int(candidate['source_song_id'])
if candidate['platform'] == PLATFORM_QQ:
qq_records_by_song.setdefault(sid, []).append(candidate)
elif candidate['platform'] == PLATFORM_KUGOU:
kugou_records_by_song.setdefault(sid, []).append(candidate)
elif candidate['platform'] == PLATFORM_NETEASE:
netease_records_by_song.setdefault(sid, []).append(candidate)
# ── 批量预取 spider 数据(整批一次 IN 查询,不逐首调用)──────────
qq_mids = list({pr['platform_unique_key'] for pr in platform_records if pr['platform'] == PLATFORM_QQ})
kugou_ids = list({int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_KUGOU})
netease_ids = list({int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_NETEASE})
qq_mids = list({
pr['platform_unique_key'] for pr in platform_records if pr['platform'] == PLATFORM_QQ
} | {
pr['platform_unique_key'] for records in qq_records_by_song.values() for pr in records
})
kugou_ids = list({
int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_KUGOU
} | {
int(pr['platform_unique_key']) for records in kugou_records_by_song.values() for pr in records
})
netease_ids = list({
int(pr['platform_unique_key']) for pr in platform_records if pr['platform'] == PLATFORM_NETEASE
} | {
int(pr['platform_unique_key']) for records in netease_records_by_song.values() for pr in records
})
qq_songs_map = fetch_qq_songs(spider_conn, qq_mids) if qq_mids else {}
kugou_songs_map = fetch_kugou_songs(spider_conn, kugou_ids) if kugou_ids else {}
......@@ -1398,6 +1693,7 @@ def run(
}
prepare_inputs = []
rejected_pending: list[dict] = []
for pending in pending_records:
src_id = int(pending['song_id'])
platform = str(pending['platform'])
......@@ -1411,9 +1707,37 @@ def run(
src_id, pending['record_id'], platform,
)
continue
if platform == PLATFORM_QQ:
selected = _select_qq_import_record(
hk_row, pr, qq_records_by_song.get(src_id, []), qq_songs_map,
)
elif platform == PLATFORM_KUGOU:
selected = _select_kugou_import_record(
hk_row, pr, kugou_records_by_song.get(src_id, []), kugou_songs_map,
)
elif platform == PLATFORM_NETEASE:
selected = _select_netease_import_record(
hk_row, pr, netease_records_by_song.get(src_id, []), netease_songs_map,
)
else:
selected = pr
if selected is None:
log.warning(
'Skip import without usable lyric/cover candidate: platform=%s source_song_id=%s record_id=%s',
platform, src_id, pending['record_id'],
)
rejected_pending.append(pending)
continue
if selected['record_id'] != pr['record_id']:
log.info(
'Replace import record: platform=%s source_song_id=%s old_record_id=%s new_record_id=%s',
platform, src_id, pr['record_id'], selected['record_id'],
)
pr = selected
prepare_inputs.append((pending, hk_row, pr))
payloads = []
asset_failed_pending: list[dict] = []
max_workers = min(MAX_IMPORT_WORKERS, len(prepare_inputs)) if prepare_inputs else 0
if max_workers:
with ThreadPoolExecutor(max_workers=max_workers) as executor:
......@@ -1433,16 +1757,48 @@ def run(
total_ok += 1
else:
total_err += 1
asset_failed_pending.append(pending)
log.warning(
"Skip song because payload is empty: %s platform %s song_id=%s",
hk_row.get('name'), pr['platform'], pending['song_id'],
)
except RequiredAssetTransferError as e:
total_err += 1
asset_failed_pending.append(pending)
log.warning(
"Skip song because required asset transfer failed: %s platform %s: %s",
hk_row.get('name'), pr['platform'], e,
)
except Exception as e:
total_err += 1
log.error("Error preparing song %s platform %s: %s",
hk_row.get('name'), pr['platform'], e)
asset_failed_pending.append(pending)
log.error(
"Skip song because unexpected error: %s platform %s: %s",
hk_row.get('name'), pr['platform'], e,
)
discarded_pending = rejected_pending + asset_failed_pending
if payloads:
with pg_conn.cursor() as pg_cur:
_write_import_payloads(pg_cur, payloads)
for pending in discarded_pending:
mark_yinyan_record_resource_failed(
pg_cur, int(pending['song_id']), int(pending['record_id']), str(pending['platform']),
)
pg_conn.commit()
if discarded_pending:
mark_hk_songs_deleted(hk_conn, [int(pending['song_id']) for pending in discarded_pending])
hk_conn.commit()
imported.extend(payload['result'] for payload in payloads)
elif discarded_pending:
with pg_conn.cursor() as pg_cur:
for pending in discarded_pending:
mark_yinyan_record_resource_failed(
pg_cur, int(pending['song_id']), int(pending['record_id']), str(pending['platform']),
)
pg_conn.commit()
mark_hk_songs_deleted(hk_conn, [int(pending['song_id']) for pending in discarded_pending])
hk_conn.commit()
else:
log.error("No import payloads were prepared in this batch; stopping to avoid retry loop")
break
......
......@@ -3,6 +3,7 @@ import uuid
_SINGER_INDEX_VALUES = set('ABCDEFGHIJKLMNOPQRSTUVWXYZ#')
_SINGER_SEX_VALUES = {'M', 'F', 'C', 'U'}
_SINGER_AREA_VALUES = {'华语', '欧美', '韩国', '日本', '其他'}
RESOURCE_TRANSFER_FAILED = 'resource_transfer_failed'
def _singer_index(value) -> str:
......@@ -49,10 +50,11 @@ def fetch_pending_yinyan_song_records(cur, limit: int) -> list[dict]:
FROM yinyan_song_records
WHERE is_yinyan_push = FALSE
AND platform IS NOT NULL
AND COALESCE(recording_id, '') != %s
ORDER BY song_id
LIMIT %s
""",
(limit,),
(RESOURCE_TRANSFER_FAILED, limit),
)
rows = cur.fetchall()
return [{'song_id': row[0], 'record_id': row[1], 'platform': row[2]} for row in rows]
......@@ -135,6 +137,90 @@ def fetch_qq_songs_missing_singers(cur, limit: int) -> list[dict]:
return [{'id': str(row[0]), 'platform_song_id': row[1], 'mid': row[2]} for row in rows]
def fetch_qq_songs_with_invalid_covers(cur, invalid_cover: str, limit: int) -> list[dict]:
"""返回已导入、但仍使用 QQ 无专辑占位封面的歌曲及其源歌曲状态。"""
cur.execute(
"""
SELECT qs.id, qs.platform_song_id, qs.mid, qs.name,
ysr.song_id, ysr.record_id
FROM crawler_qqmusic_songs qs
JOIN yinyan_song_records ysr
ON ysr.platform = '1' AND ysr.platform_song_id = qs.platform_song_id
WHERE qs.cover = %s
AND qs.album_id IS NULL
AND ysr.is_yinyan_push = TRUE
ORDER BY ysr.song_id
LIMIT %s
""",
(invalid_cover, limit),
)
return [
{
'id': str(row[0]), 'platform_song_id': int(row[1]), 'mid': row[2],
'name': row[3], 'song_id': int(row[4]), 'record_id': int(row[5]),
}
for row in cur.fetchall()
]
def replace_qq_song(cur, old_platform_song_id: int, song: dict) -> None:
"""用同一源歌曲下的另一条 QQ 录音更新既有 crawler 歌曲。"""
cur.execute(
"""
UPDATE crawler_qqmusic_songs
SET platform_song_id = %s, mid = %s, album_id = %s, cover = %s,
title = %s, name = %s, duration = %s, lyric = %s,
composer_name = %s, lyricist_name = %s, url = %s, audio_md5 = %s,
lyric_url = %s, platform_index_url = %s, published_at = %s,
album = %s::json, singers = %s::jsonb, version = %s,
provider_name = %s, crawler_source_data = %s::json, updated_at = NOW()
WHERE platform_song_id = %s
""",
(
song['platform_song_id'], song['mid'], song.get('album_id'), song.get('cover', ''),
song.get('title', ''), song.get('name', ''), song.get('duration', 0) or 0,
song.get('lyric'), song.get('composer_name'), song.get('lyricist_name'),
song.get('url', ''), song.get('audio_md5'), song.get('lyric_url'),
song.get('platform_index_url'), song.get('published_at'), song.get('album_json'),
song.get('singers_json', '[]'), song.get('version'), song.get('provider_name'),
song.get('crawler_source_data'), old_platform_song_id,
),
)
def replace_yinyan_song_record(cur, song_id: int, old_platform_song_id: int, record_id: int, platform_song_id: int) -> None:
cur.execute(
"""
UPDATE yinyan_song_records
SET record_id = %s, platform_song_id = %s, platform = '1'
WHERE song_id = %s AND platform = '1' AND platform_song_id = %s
""",
(record_id, platform_song_id, song_id, old_platform_song_id),
)
def delete_qq_song_and_yinyan_record(cur, song_id: int, platform_song_id: int, song_uuid: str) -> None:
"""移除无法找到替代录音的 crawler 数据及其导入状态。"""
cur.execute("DELETE FROM crawler_qqmusic_singer_songs WHERE song_id = %s", (song_uuid,))
cur.execute("DELETE FROM crawler_qqmusic_songs WHERE platform_song_id = %s", (platform_song_id,))
cur.execute(
"DELETE FROM yinyan_song_records WHERE song_id = %s AND platform = '1' AND platform_song_id = %s",
(song_id, platform_song_id),
)
def mark_yinyan_record_resource_failed(cur, song_id: int, record_id: int, platform: str) -> None:
"""保留未导入状态,并标记资源转存失败,避免后续批次重复处理。"""
cur.execute(
"""
UPDATE yinyan_song_records
SET is_yinyan_push = FALSE, recording_id = %s
WHERE song_id = %s AND record_id = %s AND platform = %s AND is_yinyan_push = FALSE
""",
(RESOURCE_TRANSFER_FAILED, song_id, record_id, platform),
)
def fetch_kugou_songs_missing_singers(cur, limit: int) -> list[dict]:
cur.execute(
"""
......
#!/usr/bin/env python3
import argparse
from etl_to_crawler.config import PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, PLATFORMS
from etl_to_crawler.runner import backfill_yinyan_record_platforms, initialize_yinyan_song_records, run, backfill_empty_singers
from etl_to_crawler.runner import backfill_yinyan_record_platforms, initialize_yinyan_song_records, run, backfill_empty_singers, backfill_qq_invalid_covers
PLATFORM_MAP = {
'qq': PLATFORM_QQ,
......@@ -22,6 +22,8 @@ if __name__ == '__main__':
help='只回填 yinyan_song_records 中为空的 platform 平台代码')
parser.add_argument('--backfill-empty-singers', action='store_true',
help='回填三平台 singers 为空的歌曲记录')
parser.add_argument('--backfill-qq-invalid-covers', action='store_true',
help='替换 QQ 无专辑占位封面的录音;无候选时删除 crawler 记录并软删 hk_songs_test')
args = parser.parse_args()
if args.platform == 'all':
......@@ -36,5 +38,7 @@ if __name__ == '__main__':
initialize_yinyan_song_records(platforms, max_batches=args.max_batches)
elif args.backfill_empty_singers:
backfill_empty_singers(platforms, max_batches=args.max_batches)
elif args.backfill_qq_invalid_covers:
backfill_qq_invalid_covers(max_batches=args.max_batches)
else:
run(platforms, max_batches=args.max_batches)
......
from unittest.mock import MagicMock, patch
import math
import requests
from etl_to_crawler.oss import transfer_url, transfer_url_with_md5
ARCHIVE_URL = "https://archive-dev.oss-cn-beijing.aliyuncs.com/some/path.mp3"
......@@ -75,3 +76,33 @@ def test_transfer_url_with_md5_hashes_downloaded_content_before_upload():
bucket.put_object.assert_called_once_with("crawler/qq/audio/abc.mp3", fake_content)
assert result == f"{BASE_URL}/crawler/qq/audio/abc.mp3"
assert audio_md5 == "04d43544b267629d9089eaed3b847a99"
def test_audio_transfer_retries_timeout_then_uploads():
bucket = MagicMock()
response = MagicMock(content=b'audio_bytes')
response.raise_for_status.return_value = None
with patch('etl_to_crawler.oss._http_get', side_effect=[requests.ReadTimeout('slow'), response]) as mock_get:
with patch('etl_to_crawler.oss.time.sleep') as mock_sleep:
result, _ = transfer_url_with_md5(OTHER_URL, 'crawler/qq/audio/abc.mp3', bucket, BASE_URL)
assert mock_get.call_count == 2
mock_sleep.assert_called_once()
assert result == f"{BASE_URL}/crawler/qq/audio/abc.mp3"
def test_transfer_does_not_retry_not_found_response():
bucket = MagicMock()
response = MagicMock(status_code=404)
error = requests.HTTPError(response=response)
with patch('etl_to_crawler.oss._http_get', side_effect=error) as mock_get:
with patch('etl_to_crawler.oss.time.sleep') as mock_sleep:
try:
transfer_url_with_md5(OTHER_URL, 'crawler/qq/audio/abc.mp3', bucket, BASE_URL)
except requests.HTTPError:
pass
else:
raise AssertionError('expected HTTPError')
mock_get.assert_called_once()
mock_sleep.assert_not_called()
......
from unittest.mock import MagicMock
import time
import pytest
from etl_to_crawler import runner
......@@ -35,6 +36,12 @@ class _Connection:
def close(self):
return None
def cursor(self):
return _Cursor()
def commit(self):
return None
def test_run_io_tasks_returns_named_results():
def slow(value):
......@@ -52,6 +59,86 @@ def test_run_io_tasks_returns_named_results():
}
def test_safe_audio_transfer_returns_empty_url_after_final_failure(monkeypatch):
monkeypatch.setattr(runner, 'transfer_url_with_md5', MagicMock(side_effect=TimeoutError('timeout')))
assert runner._safe_transfer_audio('https://source.example/audio.mp3', 'audio.mp3', object(), 'https://archive.example') == ('', '')
def test_safe_transfer_returns_empty_url_after_final_failure(monkeypatch):
monkeypatch.setattr(runner, 'transfer_url', MagicMock(side_effect=TimeoutError('timeout')))
assert runner._safe_transfer('https://source.example/cover.jpg', 'cover.jpg', object(), 'https://archive.example') == ''
def test_required_asset_validation_rejects_partial_transfer():
with pytest.raises(runner.RequiredAssetTransferError, match='audio, lyric'):
runner._require_primary_assets({
'audio': ('', ''),
'cover': 'https://archive.example/cover.jpg',
'lyric': '',
})
def test_lyric_upload_does_not_fall_back_to_external_source_url(monkeypatch):
monkeypatch.setattr(runner, 'upload_plain_lyric_to_bucket', lambda *args: '')
assert runner._safe_upload_lyric(
'qq', 'mid', '', 'https://source.example/lyric.txt', object(), 'https://archive.example',
) == ''
def test_qq_cover_source_ignores_missing_album_placeholder_from_hk():
assert runner._qq_cover_source(
{'cover_url': runner.QQ_MISSING_ALBUM_COVER},
{'cover': 'https://example.com/valid-qq-cover.jpg'},
) == 'https://example.com/valid-qq-cover.jpg'
def test_qq_cover_source_keeps_valid_hk_cover():
assert runner._qq_cover_source(
{'cover_url': 'https://example.com/hk-cover.jpg'},
{'cover': 'https://example.com/qq-cover.jpg'},
) == 'https://example.com/hk-cover.jpg'
def test_select_qq_import_record_replaces_missing_album_cover_with_same_source_candidate():
current = {'record_id': 10, 'platform_unique_key': 'bad-mid'}
candidate = {'record_id': 20, 'platform_unique_key': 'good-mid'}
selected = runner._select_qq_import_record(
{'cover_url': runner.QQ_MISSING_ALBUM_COVER},
current,
[current, candidate],
{
'bad-mid': {'album_id': 0, 'cover': runner.QQ_MISSING_ALBUM_COVER, 'lyric': '歌词'},
'good-mid': {'album_id': 100, 'cover': 'https://example.com/good.jpg', 'lyric': '歌词'},
},
)
assert selected == candidate
def test_select_qq_import_record_replaces_missing_lyric_with_same_source_candidate():
current = {'record_id': 10, 'platform_unique_key': 'without-lyric'}
candidate = {'record_id': 20, 'platform_unique_key': 'with-lyric'}
selected = runner._select_qq_import_record(
{'cover_url': 'https://example.com/hk-cover.jpg'},
current,
[current, candidate],
{
'without-lyric': {'album_id': 100, 'cover': 'https://example.com/current.jpg', 'lyric': ''},
'with-lyric': {'album_id': 100, 'cover': 'https://example.com/candidate.jpg', 'lyric': '有效歌词'},
},
)
assert selected == candidate
def test_select_qq_import_record_returns_none_when_no_candidate_satisfies_missing_fields():
current = {'record_id': 10, 'platform_unique_key': 'bad-mid'}
assert runner._select_qq_import_record(
{'cover_url': runner.QQ_MISSING_ALBUM_COVER},
current,
[current],
{'bad-mid': {'album_id': 0, 'cover': runner.QQ_MISSING_ALBUM_COVER, 'lyric': ''}},
) is None
def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
pg_conn = _PgConnection()
prepare = MagicMock(return_value={
......@@ -95,10 +182,22 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
'pub_time': '2021-01-01',
}]
monkeypatch.setattr(runner, 'fetch_platform_records', lambda conn, song_ids: platform_records)
monkeypatch.setattr(runner, 'fetch_all_platform_records', MagicMock())
monkeypatch.setattr(runner, 'fetch_all_platform_records', lambda conn, song_ids: [
{
'source_song_id': 10,
'record_id': 200,
'platform': '2',
'platform_unique_key': '200',
'platform_mid': 'kg-hash',
'album_audio_id': None,
'is_main_version': 1,
'is_high': 0,
'pub_time': '2021-01-01',
},
])
monkeypatch.setattr(runner, 'fetch_platform_records_by_record_ids', lambda conn, record_ids: platform_records)
monkeypatch.setattr(runner, 'fetch_kugou_songs', lambda conn, song_ids: {
200: {'id': 200},
200: {'id': 200, 'lyric': '歌词', 'cover': 'https://example.com/cover.jpg'},
})
monkeypatch.setattr(runner, 'fetch_kugou_singers', lambda conn, song_ids: {})
monkeypatch.setattr(runner, '_prepare_import_payload', prepare)
......@@ -108,7 +207,6 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
prepare.assert_called_once()
write_payloads.assert_called_once()
runner.fetch_all_platform_records.assert_not_called()
assert pg_conn.commits == 1
......@@ -213,6 +311,7 @@ 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, '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))
......@@ -275,6 +374,7 @@ 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, '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))
......@@ -317,6 +417,7 @@ 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, '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))
......
......@@ -15,6 +15,9 @@ from etl_to_crawler.writer import (
upsert_qq_songs,
upsert_qq_singer_songs,
upsert_yinyan_song_records,
replace_qq_song,
delete_qq_song_and_yinyan_record,
mark_yinyan_record_resource_failed,
)
......@@ -31,6 +34,44 @@ def test_upsert_qq_singers_executes_insert():
assert 'ON CONFLICT' in args[0][0]
def test_replace_qq_song_updates_existing_row_by_old_platform_song_id():
cur = MagicMock()
replace_qq_song(cur, 100, {
'platform_song_id': 200, 'mid': 'new-mid', 'album_id': 20,
'cover': 'https://oss.example/cover.jpg', 'title': '标题', 'name': '歌名',
'duration': 180, 'singers_json': '[]', 'provider_name': 'yinyan',
'crawler_source_data': '{}',
})
sql, params = cur.execute.call_args[0]
assert 'UPDATE crawler_qqmusic_songs' in sql
assert 'platform_song_id = %s' in sql
assert params[-1] == 100
assert params[0] == 200
def test_delete_qq_song_and_yinyan_record_cleans_song_relations_and_state():
cur = MagicMock()
delete_qq_song_and_yinyan_record(cur, 10, 100, 'song-uuid')
assert cur.execute.call_count == 3
calls = cur.execute.call_args_list
assert 'crawler_qqmusic_singer_songs' in calls[0].args[0]
assert 'crawler_qqmusic_songs' in calls[1].args[0]
assert 'yinyan_song_records' in calls[2].args[0]
def test_mark_yinyan_record_resource_failed_keeps_unpushed_state_and_excludes_retry():
cur = MagicMock()
mark_yinyan_record_resource_failed(cur, 10, 100, '1')
sql, params = cur.execute.call_args[0]
assert 'UPDATE yinyan_song_records' in sql
assert 'is_yinyan_push = FALSE' in sql
assert 'recording_id = %s' in sql
assert params == ('resource_transfer_failed', 10, 100, '1')
def test_upsert_qq_singers_empty_does_nothing():
cur = MagicMock()
upsert_qq_singers(cur, [])
......@@ -103,7 +144,7 @@ def test_fetch_pending_yinyan_song_records_reads_platform_for_single_platform_im
sql, params = cur.execute.call_args[0]
assert 'SELECT song_id, record_id, platform' in sql
assert 'platform IS NOT NULL' in sql
assert params == (500,)
assert params == ('resource_transfer_failed', 500)
assert rows == [
{'song_id': 10, 'record_id': 100, 'platform': '1'},
{'song_id': 11, 'record_id': 101, 'platform': '2'},
......