Commit 6d8ccd04 6d8ccd044b9b30fc5cc811233db4dd41ac4e6d5a by 沈秋雨

refactor(etl_to_crawler): 重构音频和歌手数据转存及导入流程

- 新增HTTP连接池配置,提升请求性能和资源复用
- 使用requests Session统一管理HTTP连接,减少请求开销
- 实现_url清洗函数,过滤无效或非法URL
- 并行执行音频、封面、歌词及歌手头像的OSS转存任务,提升效率
- 按平台统一构建导入数据payload,简化处理流程
- 批量写入歌曲、歌手、专辑及关联关系,提升数据库操作性能
- 替换原平台单独处理函数为统一预处理及写入方法
- 取消原单线程顺序处理,改为线程池并发执行数据准备和导入
- 获取平台歌曲及歌手信息支持传入缓存参数,减少重复查询
- 优化异常处理及日志记录,提高稳定性和可维护性
- 新增歌手索引和性别等字段辅助函数,完善歌手信息处理
1 parent d3694c28
......@@ -50,3 +50,5 @@ 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'))
OSS_CONNECTION_POOL_SIZE = int(os.environ.get('OSS_CONNECTION_POOL_SIZE', str(HTTP_POOL_MAXSIZE)))
......
......@@ -2,7 +2,7 @@ import pymysql
import pymysql.cursors
import pg8000
import oss2
from .config import SOURCE_DB, HK_SONGS_DB, CRAWLER_DB, OSS_CONFIG
from .config import SOURCE_DB, HK_SONGS_DB, CRAWLER_DB, OSS_CONFIG, OSS_CONNECTION_POOL_SIZE
def get_source_conn() -> pymysql.Connection:
......@@ -26,5 +26,6 @@ def get_pg_conn() -> pg8000.Connection:
def get_oss_bucket() -> oss2.Bucket:
oss2.defaults.connection_pool_size = OSS_CONNECTION_POOL_SIZE
auth = oss2.Auth(OSS_CONFIG['access_key_id'], OSS_CONFIG['access_key_secret'])
return oss2.Bucket(auth, OSS_CONFIG['endpoint'], OSS_CONFIG['bucket_name'])
......
import requests
import oss2
import math
from urllib.parse import urlparse
from .config import OSS_CONFIG
from requests.adapters import HTTPAdapter
from .config import HTTP_POOL_MAXSIZE, OSS_CONFIG
from .utils import compute_audio_md5
_HTTP_SESSION = requests.Session()
_HTTP_ADAPTER = HTTPAdapter(pool_connections=HTTP_POOL_MAXSIZE, pool_maxsize=HTTP_POOL_MAXSIZE)
_HTTP_SESSION.mount('http://', _HTTP_ADAPTER)
_HTTP_SESSION.mount('https://', _HTTP_ADAPTER)
def _http_get(url: str, timeout: int = 30):
return _HTTP_SESSION.get(url, timeout=timeout)
def _clean_url(url) -> str:
if url is None:
return ''
if isinstance(url, float) and math.isnan(url):
return ''
text = str(url).strip()
if not text or text.lower() in {'nan', 'none', 'null'}:
return ''
parsed = urlparse(text)
if parsed.scheme not in {'http', 'https'} or not parsed.netloc:
return ''
return text
def _host(url: str | None) -> str:
return urlparse(url or '').netloc.lower()
......@@ -28,11 +53,12 @@ def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: s
若 url 为空或已在目标 bucket,直接返回原 url(不上传)。
返回新的公开访问 URL。
"""
url = _clean_url(url)
if not url:
return ''
if _is_target_oss_url(url, base_url):
return url
resp = requests.get(_download_url(url, base_url), timeout=30)
resp = _http_get(_download_url(url, base_url), timeout=30)
resp.raise_for_status()
bucket.put_object(oss_key, resp.content)
return f"{base_url.rstrip('/')}/{oss_key}"
......@@ -43,11 +69,12 @@ def transfer_url_with_md5(url: str | None, oss_key: str, bucket: oss2.Bucket, ba
将音频 URL 转存到 OSS,并基于下载到的音频字节计算 MD5。
已在目标 OSS 的 URL 无需重新下载,无法可靠计算 MD5,返回空 MD5。
"""
url = _clean_url(url)
if not url:
return '', ''
if _is_target_oss_url(url, base_url):
return url, ''
resp = requests.get(_download_url(url, base_url), timeout=30)
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)
......
import uuid
_SINGER_INDEX_VALUES = set('ABCDEFGHIJKLMNOPQRSTUVWXYZ#')
_SINGER_SEX_VALUES = {'M', 'F', 'C', 'U'}
_SINGER_AREA_VALUES = {'华语', '欧美', '韩国', '日本', '其他'}
def _singer_index(value) -> str:
text = str(value or '').strip().upper()
return text if text in _SINGER_INDEX_VALUES else '#'
def _singer_sex(value) -> str:
text = str(value or '').strip().upper()
return text if text in _SINGER_SEX_VALUES else 'U'
def _singer_area(value) -> str:
text = str(value or '').strip()
return text if text in _SINGER_AREA_VALUES else '其他'
def insert_yinyan_song_records(cur, records: list[dict]) -> None:
"""Initialize yinyan_song_records rows before crawler import."""
......@@ -114,7 +133,7 @@ def upsert_qq_singers(cur, singers: list[dict]) -> None:
"""
rows = [(
s['id'], s['mid'], s['name'], s.get('avatar', ''),
s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
_singer_sex(s.get('sex')), _singer_area(s.get('area')), _singer_index(s.get('index')),
s.get('intro'), s.get('home_url'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in singers]
......@@ -206,7 +225,7 @@ def upsert_kugou_singers(cur, singers: list[dict]) -> None:
"""
rows = [(
s['id'], s['name'], s.get('avatar', ''),
s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
_singer_sex(s.get('sex')), _singer_area(s.get('area')), _singer_index(s.get('index')),
s.get('intro'), s.get('home_url', ''),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in singers]
......@@ -295,7 +314,7 @@ def upsert_netease_singers(cur, singers: list[dict]) -> None:
"""
rows = [(
s['id'], s['name'], s.get('avatar', ''),
s.get('sex') or 'U', s.get('area') or '其他', s.get('index') or '#',
_singer_sex(s.get('sex')), _singer_area(s.get('area')), _singer_index(s.get('index')),
s.get('intro'), s.get('home_url'),
s.get('provider_name'), s.get('crawler_source_data'),
) for s in singers]
......
from unittest.mock import MagicMock, patch
import math
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"
......@@ -22,10 +23,22 @@ def test_none_url_returns_empty():
result = transfer_url(None, "any/key.mp3", bucket, BASE_URL)
assert result == ''
def test_invalid_nan_url_returns_empty_without_download():
bucket = MagicMock()
with patch('etl_to_crawler.oss._http_get') as mock_get:
assert transfer_url('nan', "any/key.mp3", bucket, BASE_URL) == ''
assert transfer_url(math.nan, "any/key.mp3", bucket, BASE_URL) == ''
assert transfer_url('not-a-url', "any/key.mp3", bucket, BASE_URL) == ''
assert transfer_url_with_md5('nan', "any/key.mp3", bucket, BASE_URL) == ('', '')
mock_get.assert_not_called()
bucket.put_object.assert_not_called()
def test_external_url_downloads_and_uploads():
bucket = MagicMock()
fake_content = b"audio_bytes"
with patch('etl_to_crawler.oss.requests.get') as mock_get:
with patch('etl_to_crawler.oss._http_get') as mock_get:
mock_get.return_value.content = fake_content
mock_get.return_value.raise_for_status = MagicMock()
result = transfer_url(OTHER_URL, "crawler/qq/audio/abc.mp3", bucket, BASE_URL)
......@@ -42,7 +55,7 @@ def test_external_url_can_rewrite_download_base_to_internal_endpoint():
'download_rewrite_from_base_url': 'https://source-bucket.oss-cn-hangzhou.aliyuncs.com',
'download_base_url': internal_source_base,
}):
with patch('etl_to_crawler.oss.requests.get') as mock_get:
with patch('etl_to_crawler.oss._http_get') as mock_get:
mock_get.return_value.content = fake_content
mock_get.return_value.raise_for_status = MagicMock()
result = transfer_url(public_source, "crawler/qq/audio/abc.mp3", bucket, BASE_URL)
......@@ -55,7 +68,7 @@ def test_external_url_can_rewrite_download_base_to_internal_endpoint():
def test_transfer_url_with_md5_hashes_downloaded_content_before_upload():
bucket = MagicMock()
fake_content = b"audio_bytes"
with patch('etl_to_crawler.oss.requests.get') as mock_get:
with patch('etl_to_crawler.oss._http_get') as mock_get:
mock_get.return_value.content = fake_content
mock_get.return_value.raise_for_status = MagicMock()
result, audio_md5 = transfer_url_with_md5(OTHER_URL, "crawler/qq/audio/abc.mp3", bucket, BASE_URL)
......
from unittest.mock import MagicMock
import time
from etl_to_crawler import runner
......@@ -35,13 +36,38 @@ class _Connection:
return None
def test_run_io_tasks_returns_named_results():
def slow(value):
time.sleep(0.01)
return value
result = runner._run_io_tasks({
'audio': lambda: slow(('audio-url', 'md5')),
'cover': lambda: slow('cover-url'),
})
assert result == {
'audio': ('audio-url', 'md5'),
'cover': 'cover-url',
}
def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
pg_conn = _PgConnection()
processors = {
'1': MagicMock(return_value={'platform': 'qq', 'platform_song_id': 100, 'mid': 'qq-mid', 'title': '歌'}),
'2': MagicMock(return_value={'platform': 'kugou', 'platform_song_id': 200, 'hash': 'kg-hash', 'title': '歌'}),
}
yinyan_writer = MagicMock()
prepare = MagicMock(return_value={
'platform': '2',
'platform_song_id': 200,
'result': {'platform': 'kugou', 'platform_song_id': 200, 'hash': 'kg-hash', 'title': '歌'},
'yinyan_record': {'song_id': 10, 'record_id': 200, 'platform': '2', 'platform_song_id': 200},
'singers': [],
'albums': [],
'songs': [],
'singer_songs': [],
'singer_albums': [],
})
write_payloads = MagicMock(return_value=[
{'platform': 'kugou', 'platform_song_id': 200, 'hash': 'kg-hash', 'title': '歌'},
])
monkeypatch.setattr(runner, 'get_hk_songs_conn', lambda: _Connection())
monkeypatch.setattr(runner, 'get_source_conn', lambda: _Connection())
......@@ -57,19 +83,7 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
'audio_url': 'https://example.com/a.mp3',
'singer': '歌手',
}})
platform_records = [
{
'source_song_id': 10,
'record_id': 100,
'platform': '1',
'platform_unique_key': 'qq-mid',
'platform_mid': '100',
'album_audio_id': None,
'is_main_version': 0,
'is_high': 1,
'pub_time': '2020-01-01',
},
{
platform_records = [{
'source_song_id': 10,
'record_id': 200,
'platform': '2',
......@@ -79,25 +93,22 @@ def test_run_imports_only_pending_yinyan_platform_record(monkeypatch):
'is_main_version': 1,
'is_high': 0,
'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_platform_records_by_record_ids', lambda conn, record_ids: platform_records)
monkeypatch.setattr(runner, '_PROCESSORS', processors)
monkeypatch.setattr(runner, 'upsert_yinyan_song_records', yinyan_writer)
monkeypatch.setattr(runner, 'fetch_kugou_songs', lambda conn, song_ids: {
200: {'id': 200},
})
monkeypatch.setattr(runner, 'fetch_kugou_singers', lambda conn, song_ids: {})
monkeypatch.setattr(runner, '_prepare_import_payload', prepare)
monkeypatch.setattr(runner, '_write_import_payloads', write_payloads)
runner.run(['1', '2'])
processors['1'].assert_not_called()
processors['2'].assert_called_once()
prepare.assert_called_once()
write_payloads.assert_called_once()
runner.fetch_all_platform_records.assert_not_called()
yinyan_writer.assert_called_once_with(pg_conn.cur, [{
'song_id': 10,
'record_id': 200,
'platform': '2',
'platform_song_id': 200,
}])
assert pg_conn.commits == 1
......@@ -295,6 +306,68 @@ def test_process_netease_builds_album_json_for_song_insert(monkeypatch):
assert inserted_songs[0]['version'] == 'DJ默涵版'
def test_process_netease_uses_prefetched_song_and_singers(monkeypatch):
pg_cur = MagicMock()
pg_cur.fetchone.return_value = ('song-uuid',)
inserted_songs = []
fetch_songs = MagicMock()
fetch_singers = MagicMock()
monkeypatch.setattr(runner, 'fetch_netease_songs', fetch_songs)
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, '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))
monkeypatch.setattr(runner, 'upsert_netease_singer_songs', lambda cur, pairs: None)
runner._process_netease(
{
'name': '词曲名',
'audio_url': 'https://example.com/audio.mp3',
'lyrics_url': 'https://example.com/lyric.lrc',
'cover_url': '',
'composer': '词曲曲作者',
'lyricist': '词曲词作者',
'issue_time': '2019-01-01',
'song_time': 120,
},
{'platform_unique_key': '300'},
spider_conn=object(),
pg_cur=pg_cur,
bucket=object(),
base_url='https://bucket.example.com',
song_data={
'id': 300,
'album_id': None,
'cover': 'https://example.com/cover.jpg',
'title': '录音标题',
'duration': 180,
'lyric': '[00:01.00]歌词',
'composer_name': '曲作者',
'lyricist_name': '词作者',
'platform_index_url': None,
'published_at': '2020-01-02',
},
singer_list=[{
'singer_id': 1,
'name': '歌手',
'avatar': '',
'sex': 'U',
'area': '其他',
'index': '#',
'intro': None,
'home_url': None,
}],
)
fetch_songs.assert_not_called()
fetch_singers.assert_not_called()
assert inserted_songs[0]['platform_song_id'] == 300
assert '"name": "歌手"' in inserted_songs[0]['singers_json']
def test_process_netease_keeps_timestamped_lyric_and_uploads_plain_lyric(monkeypatch):
pg_cur = MagicMock()
pg_cur.fetchone.return_value = ('song-uuid',)
......
......@@ -387,3 +387,42 @@ def test_upsert_netease_entities_write_provider_and_source_data():
assert 'version' in sql
assert 'provider_name, crawler_source_data' in sql
assert rows[0][-4:] == ('DJ默涵版', 'md5-300', 'yinyan', '{"id": 300}')
def test_upsert_singers_normalizes_invalid_enum_values():
cur = MagicMock()
upsert_netease_singers(cur, [{
'id': 1,
'name': '歌手',
'avatar': 'https://example.com/avatar.jpg',
'sex': '0',
'area': '未知地区',
'index': '0',
'intro': '简介',
'home_url': 'https://example.com/singer',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 1}',
}])
_, rows = cur.executemany.call_args[0]
assert rows[0][3:6] == ('U', '其他', '#')
def test_upsert_singers_preserves_valid_group_sex_enum_value():
cur = MagicMock()
upsert_qq_singers(cur, [{
'id': 1,
'mid': 'singer-mid',
'name': '组合',
'avatar': 'https://example.com/avatar.jpg',
'sex': 'C',
'area': '华语',
'index': 'z',
'intro': '简介',
'home_url': 'https://example.com/singer',
'provider_name': 'yinyan',
'crawler_source_data': '{"id": 1}',
}])
_, rows = cur.executemany.call_args[0]
assert rows[0][4:7] == ('C', '华语', 'Z')
......