Commit bf8d1723 bf8d172363896b6811dcf22f5729187e27ff6398 by 沈秋雨

chore(config): 支持切换音眼导入状态表以兼容 records2

- 在配置中添加 YINYAN_IMPORT_TABLE 环境变量,支持 records 和 records2 两张表
- 增加对 YINYAN_IMPORT_TABLE 合法值的校验,防止配置错误
- 修改写入和查询相关函数以适配动态表名
- 优化 upsert 逻辑,针对 records2 版本避免多条同 song_id 记录一同更新
- run_etl.py 中打印日志增加当前使用的导入表信息
- writer.py、runner.py 调整查询及更新逻辑以支持动态导入表切换
- oss.py 新增下载歌词文本函数,确保即使目标 OSS 包含歌词文件也重新下载内容
- runner.py 优化歌词上传流程,优先转存音眼歌词 URL,失败后回退爬虫端歌词文本
- 增加多处歌词和封面判断逻辑,支持基于音眼歌词 URL 判断歌词可用性
- 添加单元测试覆盖眼眼歌词转存优先逻辑及导入状态表切换功能验证
1 parent bbcdcae8
......@@ -13,6 +13,9 @@ TARGET_DB_PASSWORD=
TARGET_DB_NAME=
TARGET_TABLE_NAME=hk_songs
# run_etl.py 默认读取并回写的导入状态表:yinyan_song_records 或 yinyan_song_records2
YINYAN_IMPORT_TABLE=yinyan_song_records
# OSS 配置
OSS_ACCESS_KEY_ID=
OSS_ACCESS_KEY_SECRET=
......
......@@ -48,6 +48,14 @@ PLATFORM_KUGOU = '2'
PLATFORM_NETEASE = '4'
PLATFORMS = [PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE]
YINYAN_IMPORT_TABLE = os.environ.get('YINYAN_IMPORT_TABLE', 'yinyan_song_records').strip()
_ALLOWED_YINYAN_IMPORT_TABLES = {'yinyan_song_records', 'yinyan_song_records2'}
if YINYAN_IMPORT_TABLE not in _ALLOWED_YINYAN_IMPORT_TABLES:
raise ValueError(
f"YINYAN_IMPORT_TABLE must be one of {sorted(_ALLOWED_YINYAN_IMPORT_TABLES)}, "
f"got {YINYAN_IMPORT_TABLE!r}"
)
BATCH_SIZE = 1000
BACKFILL_BATCH_SIZE = int(os.environ.get('BACKFILL_BATCH_SIZE', '5000'))
HTTP_POOL_MAXSIZE = int(os.environ.get('HTTP_POOL_MAXSIZE', '128'))
......
......@@ -85,6 +85,15 @@ def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: s
return f"{base_url.rstrip('/')}/{oss_key}"
def download_text_url(url: str | None, base_url: str) -> str:
"""下载歌词等文本 URL;即使源文件已在目标 OSS,也实际读取其内容。"""
url = _clean_url(url)
if not url:
return ''
content = _download_content(url, base_url)
return content.decode('utf-8-sig')
def transfer_url_with_md5(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> tuple[str, str]:
"""
将音频 URL 转存到 OSS,并基于下载到的音频字节计算 MD5。
......
import uuid
from .config import YINYAN_IMPORT_TABLE
_SINGER_INDEX_VALUES = set('ABCDEFGHIJKLMNOPQRSTUVWXYZ#')
_SINGER_SEX_VALUES = {'M', 'F', 'C', 'U'}
_SINGER_AREA_VALUES = {'华语', '欧美', '韩国', '日本', '其他'}
# 资源转存失败哨兵值:platform_song_id 设为 -1 表示 OSS 失败,区别于 NULL(待导入)和正整数(成功)
PLATFORM_SONG_ID_RESOURCE_FAILED = -1
# platform_song_id 负数为失败码;NULL 表示待导入,正整数表示成功。
YINYAN_IMPORT_FAILURE_UNKNOWN = -1
YINYAN_IMPORT_FAILURE_NO_LYRIC = -2
YINYAN_IMPORT_FAILURE_NO_COVER = -3
YINYAN_IMPORT_FAILURE_NO_LYRIC_AND_COVER = -4
YINYAN_IMPORT_FAILURE_SOURCE_DATA_MISSING = -5
YINYAN_IMPORT_FAILURE_AUDIO_TRANSFER = -6
YINYAN_IMPORT_FAILURE_COVER_TRANSFER = -7
YINYAN_IMPORT_FAILURE_LYRIC_TRANSFER = -8
YINYAN_IMPORT_FAILURE_MULTIPLE_ASSET_TRANSFER = -9
YINYAN_IMPORT_FAILURE_UNEXPECTED = -10
YINYAN_IMPORT_FAILURE_CODES = {
YINYAN_IMPORT_FAILURE_UNKNOWN: 'unknown',
YINYAN_IMPORT_FAILURE_NO_LYRIC: 'no_usable_lyric',
YINYAN_IMPORT_FAILURE_NO_COVER: 'no_usable_cover',
YINYAN_IMPORT_FAILURE_NO_LYRIC_AND_COVER: 'no_usable_lyric_and_cover',
YINYAN_IMPORT_FAILURE_SOURCE_DATA_MISSING: 'source_or_spider_data_missing',
YINYAN_IMPORT_FAILURE_AUDIO_TRANSFER: 'audio_oss_transfer_failed',
YINYAN_IMPORT_FAILURE_COVER_TRANSFER: 'cover_oss_transfer_failed',
YINYAN_IMPORT_FAILURE_LYRIC_TRANSFER: 'lyric_oss_transfer_failed',
YINYAN_IMPORT_FAILURE_MULTIPLE_ASSET_TRANSFER: 'multiple_oss_transfers_failed',
YINYAN_IMPORT_FAILURE_UNEXPECTED: 'unexpected_error',
}
def _singer_index(value) -> str:
......@@ -90,9 +114,9 @@ def fetch_existing_yinyan_song_ids(cur) -> set[int]:
def fetch_pending_yinyan_song_records(cur, limit: int) -> list[dict]:
cur.execute(
"""
f"""
SELECT song_id, record_id, platform
FROM yinyan_song_records
FROM {YINYAN_IMPORT_TABLE}
WHERE is_yinyan_push = FALSE
AND platform IS NOT NULL
AND platform_song_id IS NULL
......@@ -142,23 +166,31 @@ def update_yinyan_record_platforms(cur, records: list[dict]) -> None:
def upsert_yinyan_song_records(cur, records: list[dict]) -> None:
"""Mark pre-initialized yinyan song-record rows as pushed to crawler.
Matches only on song_id so singer-fallback can use a different record_id than initialized.
精确定位原始待处理关系,避免 records2 同一 song_id 的多条录音被一起更新。
"""
if not records:
return
values_sql = ', '.join(['(%s::bigint, %s::bigint, %s::varchar, %s::bigint)'] * len(records))
values_sql = ', '.join(
['(%s::bigint, %s::bigint, %s::bigint, %s::varchar, %s::bigint)'] * len(records)
)
params = []
for r in records:
params.extend([r['song_id'], r['record_id'], r['platform'], r['platform_song_id']])
params.extend([
r['song_id'], r.get('source_record_id', r['record_id']), r['record_id'],
r['platform'], r['platform_song_id'],
])
record_id_assignment = '' if YINYAN_IMPORT_TABLE == 'yinyan_song_records2' else 'record_id = v.record_id,'
cur.execute(
f"""
UPDATE yinyan_song_records AS ysr
UPDATE {YINYAN_IMPORT_TABLE} AS ysr
SET platform = v.platform,
platform_song_id = v.platform_song_id,
record_id = v.record_id,
{record_id_assignment}
is_yinyan_push = TRUE
FROM (VALUES {values_sql}) AS v(song_id, record_id, platform, platform_song_id)
FROM (VALUES {values_sql}) AS v(song_id, source_record_id, record_id, platform, platform_song_id)
WHERE ysr.song_id = v.song_id
AND ysr.record_id = v.source_record_id
AND ysr.platform = v.platform
AND ysr.is_yinyan_push = FALSE
""",
tuple(params),
......@@ -254,15 +286,23 @@ def delete_qq_song_and_yinyan_record(cur, song_id: int, platform_song_id: int, s
)
def mark_yinyan_record_resource_failed(cur, song_id: int, record_id: int, platform: str) -> None:
"""保留未导入状态,并将 platform_song_id 置为哨兵值 -1,避免后续批次重复处理。"""
def mark_yinyan_record_resource_failed(
cur,
song_id: int,
record_id: int,
platform: str,
failure_code: int = YINYAN_IMPORT_FAILURE_UNKNOWN,
) -> None:
"""保留未导入状态,并以负数 platform_song_id 记录具体失败类型。"""
if failure_code not in YINYAN_IMPORT_FAILURE_CODES:
failure_code = YINYAN_IMPORT_FAILURE_UNKNOWN
cur.execute(
"""
UPDATE yinyan_song_records
f"""
UPDATE {YINYAN_IMPORT_TABLE}
SET is_yinyan_push = FALSE, platform_song_id = %s
WHERE song_id = %s AND record_id = %s AND platform = %s AND is_yinyan_push = FALSE
""",
(PLATFORM_SONG_ID_RESOURCE_FAILED, song_id, record_id, platform),
(failure_code, song_id, record_id, platform),
)
......
#!/usr/bin/env python3
import argparse
from etl_to_crawler.config import PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, PLATFORMS
from etl_to_crawler.config import PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, PLATFORMS, YINYAN_IMPORT_TABLE
from etl_to_crawler.runner import backfill_yinyan_record_platforms, initialize_yinyan_song_records, initialize_yinyan_song_records2, run, backfill_empty_singers, backfill_qq_invalid_covers
PLATFORM_MAP = {
......@@ -33,7 +33,7 @@ if __name__ == '__main__':
else:
platforms = [PLATFORM_MAP[args.platform]]
print(f"Starting ETL for platforms: {platforms}")
print(f"Starting ETL for platforms: {platforms}, import table: {YINYAN_IMPORT_TABLE}")
if args.backfill_yinyan_platforms:
backfill_yinyan_record_platforms(max_batches=args.max_batches)
elif args.init_yinyan_records:
......
from unittest.mock import MagicMock, patch
import math
import requests
from etl_to_crawler.oss import transfer_url, transfer_url_with_md5
from etl_to_crawler.oss import download_text_url, transfer_url, transfer_url_with_md5
ARCHIVE_URL = "https://archive-dev.oss-cn-beijing.aliyuncs.com/some/path.mp3"
OTHER_URL = "https://hikoon-data-platform.oss-cn-beijing.aliyuncs.com/qq-audio/abc.mp3"
......@@ -13,6 +13,15 @@ def test_already_archive_dev_returns_as_is():
assert result == ARCHIVE_URL
bucket.put_object.assert_not_called()
def test_download_text_url_reads_source_even_when_it_is_target_oss():
with patch('etl_to_crawler.oss._http_get') as mock_get:
mock_get.return_value.content = b'\xef\xbb\xbf[00:01.00]lyric'
mock_get.return_value.raise_for_status = MagicMock()
result = download_text_url(ARCHIVE_URL, BASE_URL)
assert result == '[00:01.00]lyric'
def test_empty_url_returns_empty():
bucket = MagicMock()
result = transfer_url('', "any/key.mp3", bucket, BASE_URL)
......
......@@ -78,11 +78,33 @@ def test_required_asset_validation_rejects_partial_transfer():
})
def test_lyric_upload_does_not_fall_back_to_external_source_url(monkeypatch):
monkeypatch.setattr(runner, 'upload_plain_lyric_to_bucket', lambda *args: '')
def test_lyric_upload_prefers_hk_url_and_forces_oss_transfer(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(
'qq', 'mid', 'spider lyric', 'https://source.example/lyric.txt', bucket, 'https://archive.example',
)
assert result == 'https://archive.example/crawler/lyric/qq/mid.txt'
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):
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(
'qq', 'mid', '', 'https://source.example/lyric.txt', object(), 'https://archive.example',
) == ''
'qq', 'mid', 'spider lyric', 'https://source.example/lyric.txt', object(), 'https://archive.example',
) == 'https://archive.example/crawler/lyric/qq/mid.txt'
assert upload.call_args.args[2] == 'spider lyric'
def test_qq_cover_source_ignores_missing_album_placeholder_from_hk():
......@@ -131,6 +153,25 @@ def test_select_qq_import_record_replaces_missing_lyric_with_same_source_candida
assert reason == 'ok'
def test_select_qq_import_record_accepts_hk_lyric_when_spider_lyric_is_empty():
current = {'record_id': 10, 'platform_unique_key': 'without-spider-lyric'}
selected, reason = runner._select_qq_import_record(
{'lyrics_url': 'https://archive-dev.example/lyrics/song.txt', 'cover_url': 'https://example.com/cover.jpg'},
current,
[current],
{
'without-spider-lyric': {
'album_id': 100,
'cover': 'https://example.com/cover.jpg',
'lyric': '',
},
},
)
assert selected == current
assert reason == 'ok'
def test_select_qq_import_record_returns_none_when_no_candidate_satisfies_missing_fields():
current = {'record_id': 10, 'platform_unique_key': 'bad-mid'}
result, reason = runner._select_qq_import_record(
......
from unittest.mock import MagicMock, call
import etl_to_crawler.writer as writer
from etl_to_crawler.writer import (
fetch_pending_yinyan_song_records,
fetch_yinyan_records_missing_platform,
......@@ -69,12 +70,19 @@ def test_mark_yinyan_record_resource_failed_keeps_unpushed_state_and_excludes_re
mark_yinyan_record_resource_failed(cur, 10, 100, '1')
sql, params = cur.execute.call_args[0]
assert 'UPDATE yinyan_song_records' in sql
assert f'UPDATE {writer.YINYAN_IMPORT_TABLE}' in sql
assert 'is_yinyan_push = FALSE' in sql
assert 'platform_song_id = %s' in sql
assert params == (-1, 10, 100, '1')
def test_mark_yinyan_record_resource_failed_writes_specific_failure_code():
cur = MagicMock()
mark_yinyan_record_resource_failed(cur, 10, 100, '1', writer.YINYAN_IMPORT_FAILURE_NO_LYRIC)
assert cur.execute.call_args[0][1] == (-2, 10, 100, '1')
def test_upsert_qq_singers_empty_does_nothing():
cur = MagicMock()
upsert_qq_singers(cur, [])
......@@ -97,16 +105,21 @@ def test_upsert_yinyan_song_records_marks_existing_relation_as_pushed():
])
sql, params = cur.execute.call_args[0]
assert 'UPDATE yinyan_song_records' in sql
assert 'UPDATE yinyan_song_records AS ysr' in sql
assert f'UPDATE {writer.YINYAN_IMPORT_TABLE}' in sql
assert f'UPDATE {writer.YINYAN_IMPORT_TABLE} AS ysr' in sql
assert 'SET platform = v.platform' in sql
assert 'platform_song_id = v.platform_song_id' in sql
assert 'record_id = v.record_id' in sql
if writer.YINYAN_IMPORT_TABLE == 'yinyan_song_records2':
assert 'record_id = v.record_id' not in sql
else:
assert 'record_id = v.record_id' in sql
assert 'is_yinyan_push = TRUE' in sql
assert 'FROM (VALUES (%s::bigint, %s::bigint, %s::varchar, %s::bigint), (%s::bigint, %s::bigint, %s::varchar, %s::bigint))' in sql
assert 'source_record_id' in sql
assert 'ysr.song_id = v.song_id' in sql
assert 'ysr.record_id = v.source_record_id' in sql
assert 'ysr.platform = v.platform' in sql
assert 'ysr.is_yinyan_push = FALSE' in sql
assert params == (10, 100, '1', 1000, 11, 101, '2', 1001)
assert params == (10, 100, 100, '1', 1000, 11, 101, 101, '2', 1001)
def test_insert_yinyan_song_records_initializes_unpushed_rows():
......@@ -193,6 +206,23 @@ def test_fetch_pending_yinyan_song_records_reads_platform_for_single_platform_im
]
def test_import_state_queries_can_target_records2(monkeypatch):
monkeypatch.setattr(writer, 'YINYAN_IMPORT_TABLE', 'yinyan_song_records2')
cur = MagicMock()
cur.fetchall.return_value = []
writer.fetch_pending_yinyan_song_records(cur, 100)
assert 'FROM yinyan_song_records2' in cur.execute.call_args[0][0]
writer.upsert_yinyan_song_records(cur, [
{'song_id': 10, 'record_id': 100, 'platform': '1', 'platform_song_id': 1000},
])
assert 'UPDATE yinyan_song_records2 AS ysr' in cur.execute.call_args[0][0]
writer.mark_yinyan_record_resource_failed(cur, 10, 100, '1')
assert 'UPDATE yinyan_song_records2' in cur.execute.call_args[0][0]
def test_update_yinyan_record_platforms_fills_only_null_platform_rows():
cur = MagicMock()
update_yinyan_record_platforms(cur, [
......
| **失败码** | **含义** |
| ---------- | ---------------------------- |
| -1 | 未分类/未知失败 |
| -2 | 没有可用歌词 |
| -3 | 没有可用封面 |
| -4 | 歌词和封面均不可用 |
| -5 | HK、音眼源库或spider数据缺失 |
| -6 | 音频转存OSS失败 |
| -7 | 封面转存OSS失败 |
| -8 | 歌词转存OSS失败 |
| -9 | 多项必需资源转存失败 |
| -10 | 未预期程序异常 |
\ No newline at end of file