2026-07-08-hk-songs-to-crawler-etl.md 50.6 KB

hk_songs_test → crawler_dev ETL 实现计划

For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (- [ ]) syntax for tracking.

Goal:new_music_library.hk_songs_test(87K 条,deleted=0)通过 hikoon-data 关联链补充平台信息、hikoon-data-spider 补充完整字段,OSS 转移所有媒体 URL,批量写入 crawler_dev QQ/酷狗/网易三平台共 15 张表(songs/singers/albums/singer_songs/singer_albums)。

Architecture: 单进程 Python ETL,分层模块:config → connections → reader(hk_songs_test + hikoon-data 关联)→ spider(hikoon-data-spider 详细数据)→ oss/lyric 处理 → writer(PostgreSQL 批量写入)→ runner(批次编排)。每步幂等,按唯一键 ON CONFLICT DO NOTHING,可重跑。

Tech Stack: Python 3.10、pymysql、psycopg2-binary、oss2、requests、tqdm、python-dotenv;运行于项目根目录 .venv/

Global Constraints

  • virtualenv: 项目根目录 .venv/,运行前 source .venv/bin/activate
  • 所有外部 URL(音频、封面、歌手头像)统一转到 archive-dev.oss-cn-beijing.aliyuncs.com;已在 archive-dev 的 URL 直接复用
  • 写入顺序每平台固定:singers → albums → songs → singer_songs → singer_albums
  • 所有 PostgreSQL 写入按唯一键 ON CONFLICT DO NOTHING(幂等)
  • 批量大小:每批 500 条 hk_songs_test 记录
  • 平台编号:1=QQ,2=酷狗,4=网易
  • 新增 .env 变量:CRAWLER_DB_HOSTCRAWLER_DB_PORTCRAWLER_DB_USERCRAWLER_DB_PASSWORDCRAWLER_DB_NAME(对应 crawler_dev PostgreSQL)

File Structure

etl_to_crawler/
    __init__.py        # 空文件
    config.py          # 从 .env 加载全部连接参数,暴露常量
    connections.py     # 创建并返回 MySQL / PostgreSQL / OSS 连接对象
    lyric.py           # 去除 LRC [mm:ss.xx] 时间戳
    oss.py             # 将任意 URL 转移到 archive-dev OSS,已在则跳过
    reader.py          # 读 hk_songs_test(TencentDB)+ hikoon-data 关联链
    spider.py          # 从 hikoon-data-spider 获取三平台详细数据
    writer.py          # 批量写入 PostgreSQL(5 类表 × 3 平台)
    runner.py          # 批次编排、进度日志、错误收集

tests/
    test_lyric.py
    test_oss.py
    test_spider.py
    test_writer.py

run_etl.py             # CLI 入口:python run_etl.py [--platform qq|kugou|netease|all]

Task 1: 环境准备 & 连接验证

Files:

  • Modify: requirements.txt
  • Create: etl_to_crawler/__init__.py
  • Create: etl_to_crawler/config.py
  • Create: etl_to_crawler/connections.py
  • Modify: .env(追加 CRAWLER_DB_* 变量)

Interfaces:

  • Produces:

    • config.SOURCE_DB: dict — hikoon-data MySQL 连接参数
    • config.HK_SONGS_DB: dict — new_music_library MySQL 连接参数
    • config.CRAWLER_DB: dict — crawler_dev PostgreSQL 连接参数
    • config.OSS_CONFIG: dict — archive-dev OSS 参数
    • connections.get_source_conn() -> pymysql.Connection — hikoon-data
    • connections.get_hk_songs_conn() -> pymysql.Connection — new_music_library
    • connections.get_pg_conn() -> psycopg2.connection — crawler_dev
    • connections.get_oss_bucket() -> oss2.Bucket
  • Step 1: 安装 psycopg2-binary

source .venv/bin/activate
pip install psycopg2-binary
pip freeze | grep psycopg2

期望输出:psycopg2-binary==2.9.x

  • Step 2: 追加 CRAWLER_DB 变量到 .env

.env 末尾追加(填入实际值):

CRAWLER_DB_HOST=data-center.rwlb.rds.aliyuncs.com
CRAWLER_DB_PORT=5432
CRAWLER_DB_USER=data_test
CRAWLER_DB_PASSWORD=z1n_z#JHW0
CRAWLER_DB_NAME=crawler_dev
  • Step 3: 创建 etl_to_crawler/__init__.py

(空文件)

  • Step 4: 创建 etl_to_crawler/config.py
import os
from dotenv import load_dotenv

load_dotenv()

SOURCE_DB = {
    'host': os.environ['SOURCE_DB_HOST'],
    'port': int(os.environ['SOURCE_DB_PORT']),
    'user': os.environ['SOURCE_DB_USER'],
    'password': os.environ['SOURCE_DB_PASSWORD'],
    'database': os.environ['SOURCE_DB_NAME'],
    'charset': 'utf8mb4',
}

HK_SONGS_DB = {
    'host': os.environ['TARGET_DB_HOST'],
    'port': int(os.environ['TARGET_DB_PORT']),
    'user': os.environ['TARGET_DB_USER'],
    'password': os.environ['TARGET_DB_PASSWORD'],
    'database': os.environ['TARGET_DB_NAME'],
    'charset': 'utf8mb4',
}

CRAWLER_DB = {
    'host': os.environ['CRAWLER_DB_HOST'],
    'port': int(os.environ.get('CRAWLER_DB_PORT', '5432')),
    'user': os.environ['CRAWLER_DB_USER'],
    'password': os.environ['CRAWLER_DB_PASSWORD'],
    'dbname': os.environ['CRAWLER_DB_NAME'],
    'sslmode': 'prefer',
}

OSS_CONFIG = {
    'access_key_id': os.environ['OSS_ACCESS_KEY_ID'],
    'access_key_secret': os.environ['OSS_ACCESS_KEY_SECRET'],
    'endpoint': os.environ['OSS_ENDPOINT'],
    'bucket_name': os.environ['OSS_BUCKET_NAME'],
    'base_url': f"https://{os.environ['OSS_BUCKET_NAME']}.{os.environ['OSS_ENDPOINT']}",
}

PLATFORM_QQ = '1'
PLATFORM_KUGOU = '2'
PLATFORM_NETEASE = '4'
PLATFORMS = [PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE]

BATCH_SIZE = 500
  • Step 5: 创建 etl_to_crawler/connections.py
import pymysql
import pymysql.cursors
import psycopg2
import oss2
from .config import SOURCE_DB, HK_SONGS_DB, CRAWLER_DB, OSS_CONFIG


def get_source_conn() -> pymysql.Connection:
    return pymysql.connect(**SOURCE_DB, cursorclass=pymysql.cursors.DictCursor)


def get_spider_conn() -> pymysql.Connection:
    cfg = SOURCE_DB.copy()
    cfg['database'] = 'hikoon-data-spider'
    return pymysql.connect(**cfg, cursorclass=pymysql.cursors.DictCursor)


def get_hk_songs_conn() -> pymysql.Connection:
    return pymysql.connect(**HK_SONGS_DB, cursorclass=pymysql.cursors.DictCursor)


def get_pg_conn() -> psycopg2.extensions.connection:
    return psycopg2.connect(**CRAWLER_DB)


def get_oss_bucket() -> oss2.Bucket:
    auth = oss2.Auth(OSS_CONFIG['access_key_id'], OSS_CONFIG['access_key_secret'])
    return oss2.Bucket(auth, OSS_CONFIG['endpoint'], OSS_CONFIG['bucket_name'])
  • Step 6: 验证所有连接
source .venv/bin/activate
python3 -c "
from etl_to_crawler.connections import (
    get_source_conn, get_spider_conn, get_hk_songs_conn, get_pg_conn, get_oss_bucket
)
c = get_source_conn(); c.close(); print('hikoon-data OK')
c = get_spider_conn(); c.close(); print('hikoon-data-spider OK')
c = get_hk_songs_conn(); c.close(); print('new_music_library OK')
c = get_pg_conn(); c.close(); print('crawler_dev OK')
b = get_oss_bucket(); b.get_bucket_info(); print('OSS OK')
"

期望输出:5 行全部 OK。

  • Step 7: 更新 requirements.txt
pip freeze > requirements.txt
  • Step 8: Commit
git add requirements.txt .env etl_to_crawler/
git commit -m "feat(etl): 初始化 etl_to_crawler 模块,验证三端连接"

Task 2: lyric.py — 去除 LRC 时间戳

Files:

  • Create: etl_to_crawler/lyric.py
  • Create: tests/test_lyric.py

Interfaces:

  • Produces: lyric.strip_timestamps(text: str | None) -> str

  • Step 1: 写测试

# tests/test_lyric.py
from etl_to_crawler.lyric import strip_timestamps

def test_strips_standard_timestamps():
    lrc = "[00:12.34]第一行歌词\n[01:23.45]第二行歌词"
    assert strip_timestamps(lrc) == "第一行歌词\n第二行歌词"

def test_strips_millisecond_timestamps():
    lrc = "[00:12.345]带三位毫秒"
    assert strip_timestamps(lrc) == "带三位毫秒"

def test_strips_meta_tags():
    lrc = "[ti:歌曲名]\n[ar:歌手]\n[00:01.00]歌词"
    result = strip_timestamps(lrc)
    assert "歌词" in result
    assert "[ti:" not in result

def test_empty_and_none():
    assert strip_timestamps(None) == ''
    assert strip_timestamps('') == ''

def test_plain_text_unchanged():
    assert strip_timestamps("没有时间戳的歌词") == "没有时间戳的歌词"
  • Step 2: 运行测试确认失败
source .venv/bin/activate && pytest tests/test_lyric.py -v

期望:FAILED(ImportError 或 NameError)

  • Step 3: 实现 etl_to_crawler/lyric.py
import re

_TIMESTAMP_RE = re.compile(r'\[\d{2}:\d{2}\.\d{2,3}\]')
_META_TAG_RE = re.compile(r'\[[a-zA-Z]+:[^\]]*\]')


def strip_timestamps(text: str | None) -> str:
    if not text:
        return ''
    text = _META_TAG_RE.sub('', text)
    text = _TIMESTAMP_RE.sub('', text)
    lines = [line.strip() for line in text.splitlines() if line.strip()]
    return '\n'.join(lines)
  • Step 4: 运行测试确认通过
pytest tests/test_lyric.py -v

期望:5 个测试全部 PASSED

  • Step 5: Commit
git add etl_to_crawler/lyric.py tests/test_lyric.py
git commit -m "feat(etl): 实现 LRC 时间戳去除"

Task 3: oss.py — URL 转移到 archive-dev

Files:

  • Create: etl_to_crawler/oss.py
  • Create: tests/test_oss.py

Interfaces:

  • Consumes: connections.get_oss_bucket()config.OSS_CONFIG
  • Produces: oss.transfer_url(url: str, oss_key: str, bucket: oss2.Bucket, base_url: str) -> str

  • Step 1: 写测试

# tests/test_oss.py
from unittest.mock import MagicMock, patch
from etl_to_crawler.oss import transfer_url

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"
BASE_URL = "https://archive-dev.oss-cn-beijing.aliyuncs.com"

def test_already_archive_dev_returns_as_is():
    bucket = MagicMock()
    result = transfer_url(ARCHIVE_URL, "any/key.mp3", bucket, BASE_URL)
    assert result == ARCHIVE_URL
    bucket.put_object.assert_not_called()

def test_empty_url_returns_empty():
    bucket = MagicMock()
    result = transfer_url('', "any/key.mp3", bucket, BASE_URL)
    assert result == ''
    bucket.put_object.assert_not_called()

def test_none_url_returns_empty():
    bucket = MagicMock()
    result = transfer_url(None, "any/key.mp3", bucket, BASE_URL)
    assert result == ''

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:
        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)
    bucket.put_object.assert_called_once_with("crawler/qq/audio/abc.mp3", fake_content)
    assert result == f"{BASE_URL}/crawler/qq/audio/abc.mp3"
  • Step 2: 运行测试确认失败
pytest tests/test_oss.py -v

期望:FAILED(ImportError)

  • Step 3: 实现 etl_to_crawler/oss.py
import requests
import oss2
from urllib.parse import urlparse

ARCHIVE_DEV_HOST = 'archive-dev.oss-cn-beijing.aliyuncs.com'


def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> str:
    """
    将 url 指向的文件转移到 archive-dev OSS 的 oss_key 路径。
    若 url 为空或已在 archive-dev,直接返回原 url(不上传)。
    返回新的 archive-dev URL。
    """
    if not url:
        return ''
    if ARCHIVE_DEV_HOST in url:
        return url
    resp = requests.get(url, timeout=30)
    resp.raise_for_status()
    bucket.put_object(oss_key, resp.content)
    return f"{base_url.rstrip('/')}/{oss_key}"


def build_oss_key(platform: str, category: str, filename: str) -> str:
    """
    构造 archive-dev 内的存储路径。
    platform: 'qq' | 'kugou' | 'netease'
    category: 'audio' | 'cover' | 'singer' | 'album'
    filename: 带扩展名的文件名
    """
    return f"crawler/{platform}/{category}/{filename}"
  • Step 4: 运行测试确认通过
pytest tests/test_oss.py -v

期望:4 个测试全部 PASSED

  • Step 5: Commit
git add etl_to_crawler/oss.py tests/test_oss.py
git commit -m "feat(etl): 实现 OSS URL 转移工具"

Task 4: reader.py — 读取源数据

Files:

  • Create: etl_to_crawler/reader.py

Interfaces:

  • Consumes: connections.get_hk_songs_conn()connections.get_source_conn()
  • Produces:
    • reader.iter_hk_songs_batches(conn, batch_size: int) -> Iterator[list[dict]] — 每批返回 hk_songs_test 记录
    • reader.fetch_platform_records(source_conn, song_ids: list[int]) -> list[dict] — 返回每条的平台信息

每条 platform record 格式:

  {
    'source_song_id': int,          # hk_song_platform.id
    'platform': str,                # '1'|'2'|'4'
    'platform_unique_key': str,     # QQ:mid / 酷狗:数字ID / 网易:数字ID
    'platform_mid': str,            # QQ:数字ID / 酷狗:hash
    'album_audio_id': str,
  }
  • Step 1: 实现 etl_to_crawler/reader.py
from typing import Iterator
import pymysql
from .config import PLATFORMS


_HK_SONGS_QUERY = """
SELECT id, name, lyricist, composer, audio_url, lyrics_url,
       cover_url, singer, issue_time, source_song_id
FROM hk_songs_test
WHERE deleted = '0'
  AND name IS NOT NULL AND name != ''
  AND audio_url IS NOT NULL AND audio_url != ''
  AND singer IS NOT NULL AND singer != ''
ORDER BY id
LIMIT %s OFFSET %s
"""

_PLATFORM_QUERY = """
SELECT
    sar.song_id AS source_song_id,
    mr.platform,
    mr.platform_unique_key,
    mr.platform_mid,
    mr.album_audio_id,
    sar.is_main_version
FROM hk_song_and_record sar
JOIN hk_music_record mr ON mr.id = sar.record_id
WHERE sar.song_id IN ({placeholders})
  AND mr.platform IN ('1','2','4')
  AND mr.platform_unique_key IS NOT NULL
  AND mr.platform_unique_key != ''
  AND mr.deleted = 0
ORDER BY sar.song_id, mr.platform, sar.is_main_version DESC
"""


def iter_hk_songs_batches(conn: pymysql.Connection, batch_size: int) -> Iterator[list[dict]]:
    offset = 0
    with conn.cursor() as cur:
        while True:
            cur.execute(_HK_SONGS_QUERY, (batch_size, offset))
            rows = cur.fetchall()
            if not rows:
                break
            yield rows
            if len(rows) < batch_size:
                break
            offset += batch_size


def fetch_platform_records(source_conn: pymysql.Connection, song_ids: list[int]) -> list[dict]:
    if not song_ids:
        return []
    placeholders = ','.join(['%s'] * len(song_ids))
    query = _PLATFORM_QUERY.format(placeholders=placeholders)
    with source_conn.cursor() as cur:
        cur.execute(query, song_ids)
        rows = cur.fetchall()

    # 每个 (source_song_id, platform) 只保留 is_main_version 最高的一条
    seen = {}
    for row in rows:
        key = (row['source_song_id'], row['platform'])
        if key not in seen:
            seen[key] = row
    return list(seen.values())
  • Step 2: 手动验证
source .venv/bin/activate
python3 -c "
from etl_to_crawler.connections import get_hk_songs_conn, get_source_conn
from etl_to_crawler.reader import iter_hk_songs_batches, fetch_platform_records

hk = get_hk_songs_conn()
src = get_source_conn()
batches = iter_hk_songs_batches(hk, 10)
first_batch = next(batches)
print('hk_songs sample:', first_batch[0]['name'], first_batch[0]['source_song_id'])

song_ids = [int(r['source_song_id']) for r in first_batch]
recs = fetch_platform_records(src, song_ids)
print('platform records:', len(recs))
for r in recs[:3]:
    print(' ', r['platform'], r['platform_unique_key'][:20])
hk.close(); src.close()
"

期望:打印出歌曲名、source_song_id,及多条带平台号的记录。

  • Step 3: Commit
git add etl_to_crawler/reader.py
git commit -m "feat(etl): 实现 hk_songs_test 读取与 hikoon-data 关联"

Task 5: spider.py — 从 hikoon-data-spider 获取三平台详细数据

Files:

  • Create: etl_to_crawler/spider.py
  • Create: tests/test_spider.py

Interfaces:

  • Consumes: connections.get_spider_conn()
  • Produces:

    • spider.fetch_qq_songs(conn, mids: list[str]) -> dict[str, dict] — keyed by mid
    • spider.fetch_qq_singers(conn, song_ids: list[int]) -> dict[int, list[dict]] — keyed by song_id
    • spider.fetch_kugou_songs(conn, song_ids: list[int]) -> dict[int, dict] — keyed by id
    • spider.fetch_kugou_singers(conn, song_ids: list[int]) -> dict[int, list[dict]]
    • spider.fetch_netease_songs(conn, song_ids: list[int]) -> dict[int, dict]
    • spider.fetch_netease_singers(conn, song_ids: list[int]) -> dict[int, list[dict]]
  • Step 1: 写测试

# tests/test_spider.py
from unittest.mock import MagicMock, patch
from etl_to_crawler.spider import fetch_qq_songs, fetch_qq_singers

def _make_cursor(rows):
    cur = MagicMock()
    cur.__enter__ = MagicMock(return_value=cur)
    cur.__exit__ = MagicMock(return_value=False)
    cur.fetchall.return_value = rows
    return cur

def test_fetch_qq_songs_keys_by_mid():
    conn = MagicMock()
    conn.cursor.return_value = _make_cursor([
        {'mid': 'abc123', 'id': 1, 'title': '歌曲A', 'duration': 200,
         'lyric': None, 'composer_name': None, 'lyricist_name': None,
         'platform_index_url': None, 'published_at': None,
         'cover': '', 'album_id': None,
         'album_mid': None, 'album_title': None, 'album_cover': None,
         'album_intro': None, 'album_type': None, 'company_id': 0,
         'company': None, 'is_owner': 0, 'album_published_at': None}
    ])
    result = fetch_qq_songs(conn, ['abc123'])
    assert 'abc123' in result
    assert result['abc123']['title'] == '歌曲A'

def test_fetch_qq_songs_empty_list():
    conn = MagicMock()
    result = fetch_qq_songs(conn, [])
    assert result == {}
    conn.cursor.assert_not_called()

def test_fetch_qq_singers_keys_by_song_id():
    conn = MagicMock()
    conn.cursor.return_value = _make_cursor([
        {'song_id': 1, 'singer_id': 10, 'mid': 'sg1', 'name': '歌手A',
         'avatar': '', 'sex': 'M', 'area': '华语', 'index': 'G',
         'intro': None, 'home_url': None}
    ])
    result = fetch_qq_singers(conn, [1])
    assert 1 in result
    assert result[1][0]['name'] == '歌手A'
  • Step 2: 运行测试确认失败
pytest tests/test_spider.py -v

期望:FAILED(ImportError)

  • Step 3: 实现 etl_to_crawler/spider.py
import pymysql


def _ids_query(query_template: str, ids: list, conn: pymysql.Connection) -> list[dict]:
    if not ids:
        return []
    placeholders = ','.join(['%s'] * len(ids))
    query = query_template.format(placeholders=placeholders)
    with conn.cursor() as cur:
        cur.execute(query, ids)
        return cur.fetchall()


# ─── QQ Music ────────────────────────────────────────────────────────────────

_QQ_SONGS_SQL = """
SELECT s.id, s.mid, s.album_id, s.cover, s.title, s.duration,
       s.lyric, s.composer_name, s.lyricist_name, s.platform_index_url, s.published_at,
       a.mid AS album_mid, a.cover AS album_cover, a.title AS album_title,
       a.intro AS album_intro, a.type AS album_type, a.company_id,
       a.company, a.is_owner, a.published_at AS album_published_at
FROM media_tencent_songs s
LEFT JOIN media_tencent_albums a ON a.id = s.album_id
WHERE s.mid IN ({placeholders})
"""

_QQ_SINGERS_SQL = """
SELECT shs.song_id, shs.singer_id, sg.mid, sg.name, sg.avatar,
       sg.sex, sg.area, sg.`index`, sg.intro, sg.home_url
FROM media_tencent_singer_has_songs shs
JOIN media_tencent_singers sg ON sg.id = shs.singer_id
WHERE shs.song_id IN ({placeholders})
"""


def fetch_qq_songs(conn: pymysql.Connection, mids: list[str]) -> dict[str, dict]:
    rows = _ids_query(_QQ_SONGS_SQL, mids, conn)
    return {row['mid']: row for row in rows}


def fetch_qq_singers(conn: pymysql.Connection, song_ids: list[int]) -> dict[int, list[dict]]:
    rows = _ids_query(_QQ_SINGERS_SQL, song_ids, conn)
    result: dict[int, list] = {}
    for row in rows:
        result.setdefault(row['song_id'], []).append(row)
    return result


# ─── Kugou ───────────────────────────────────────────────────────────────────

_KUGOU_SONGS_SQL = """
SELECT s.id, s.hid, s.album_audio_id, s.album_id, s.cover, s.title, s.duration,
       s.lyric, s.composer_name, s.lyricist_name, s.platform_index_url, s.published_at,
       a.cover AS album_cover, a.title AS album_title, a.intro AS album_intro,
       a.type AS album_type, a.company_id, a.company, a.is_owner,
       a.published_at AS album_published_at
FROM media_ku_gou_songs s
LEFT JOIN media_ku_gou_albums a ON a.id = s.album_id
WHERE s.id IN ({placeholders})
"""

_KUGOU_SINGERS_SQL = """
SELECT shs.song_id, shs.singer_id, sg.name, sg.avatar,
       sg.sex, sg.area, sg.`index`, sg.intro, sg.home_url
FROM media_ku_gou_singer_has_songs shs
JOIN media_ku_gou_singers sg ON sg.id = shs.singer_id
WHERE shs.song_id IN ({placeholders})
"""


def fetch_kugou_songs(conn: pymysql.Connection, song_ids: list[int]) -> dict[int, dict]:
    rows = _ids_query(_KUGOU_SONGS_SQL, song_ids, conn)
    return {row['id']: row for row in rows}


def fetch_kugou_singers(conn: pymysql.Connection, song_ids: list[int]) -> dict[int, list[dict]]:
    rows = _ids_query(_KUGOU_SINGERS_SQL, song_ids, conn)
    result: dict[int, list] = {}
    for row in rows:
        result.setdefault(row['song_id'], []).append(row)
    return result


# ─── Netease ─────────────────────────────────────────────────────────────────

_NETEASE_SONGS_SQL = """
SELECT s.id, s.album_id, s.cover, s.title, s.duration,
       s.lyric, s.composer_name, s.lyricist_name, s.platform_index_url, s.published_at,
       a.cover AS album_cover, a.title AS album_title, a.intro AS album_intro,
       a.type AS album_type, a.company_id, a.company, a.is_owner,
       a.published_at AS album_published_at
FROM media_netease_songs s
LEFT JOIN media_netease_albums a ON a.id = s.album_id
WHERE s.id IN ({placeholders})
"""

_NETEASE_SINGERS_SQL = """
SELECT shs.song_id, shs.singer_id, sg.name, sg.avatar,
       sg.sex, sg.area, sg.`index`, sg.intro, sg.home_url
FROM media_netease_singer_has_songs shs
JOIN media_netease_singers sg ON sg.id = shs.singer_id
WHERE shs.song_id IN ({placeholders})
"""


def fetch_netease_songs(conn: pymysql.Connection, song_ids: list[int]) -> dict[int, dict]:
    rows = _ids_query(_NETEASE_SONGS_SQL, song_ids, conn)
    return {row['id']: row for row in rows}


def fetch_netease_singers(conn: pymysql.Connection, song_ids: list[int]) -> dict[int, list[dict]]:
    rows = _ids_query(_NETEASE_SINGERS_SQL, song_ids, conn)
    result: dict[int, list] = {}
    for row in rows:
        result.setdefault(row['song_id'], []).append(row)
    return result
  • Step 4: 运行测试确认通过
pytest tests/test_spider.py -v

期望:3 个测试全部 PASSED

  • Step 5: Commit
git add etl_to_crawler/spider.py tests/test_spider.py
git commit -m "feat(etl): 实现 hikoon-data-spider 三平台数据获取"

Task 6: writer.py — 写入 PostgreSQL

Files:

  • Create: etl_to_crawler/writer.py
  • Create: tests/test_writer.py

Interfaces:

  • Consumes: connections.get_pg_conn()
  • Produces:
    • writer.upsert_qq_singers(pg_cur, singers: list[dict]) -> None
    • writer.upsert_qq_albums(pg_cur, albums: list[dict]) -> None
    • writer.upsert_qq_songs(pg_cur, songs: list[dict]) -> None
    • writer.upsert_qq_singer_songs(pg_cur, pairs: list[tuple[int,str]]) -> None
    • writer.upsert_qq_singer_albums(pg_cur, pairs: list[tuple[int,int]]) -> None
    • (kugou / netease 同名函数,前缀不同)

singers 格式:{'id':int, 'mid':str, 'name':str, 'avatar':str, 'sex':str, 'area':str, 'index':str, 'intro':str|None, 'home_url':str|None}

songs 格式:{'song_uuid':str, 'platform_song_id':int, 'mid':str, 'album_id':int|None, 'cover':str, 'title':str, 'name':str, 'duration':int, 'lyric':str|None, 'composer_name':str|None, 'lyricist_name':str|None, 'url':str, 'lyric_url':str|None, 'platform_index_url':str|None, 'published_at':date|None, 'singers_json':str}

  • Step 1: 写测试
# tests/test_writer.py
from unittest.mock import MagicMock, call
from etl_to_crawler.writer import upsert_qq_singers, upsert_qq_singer_songs

def test_upsert_qq_singers_executes_insert():
    cur = MagicMock()
    singers = [{
        'id': 1, 'mid': 'abc', 'name': '歌手', 'avatar': 'https://archive-dev.oss-cn-beijing.aliyuncs.com/a.jpg',
        'sex': 'M', 'area': '华语', 'index': 'G', 'intro': None, 'home_url': None
    }]
    upsert_qq_singers(cur, singers)
    assert cur.executemany.called
    args = cur.executemany.call_args
    assert 'crawler_qqmusic_singers' in args[0][0]
    assert 'ON CONFLICT' in args[0][0]

def test_upsert_qq_singers_empty_does_nothing():
    cur = MagicMock()
    upsert_qq_singers(cur, [])
    cur.executemany.assert_not_called()

def test_upsert_qq_singer_songs():
    cur = MagicMock()
    upsert_qq_singer_songs(cur, [(1, 'uuid-abc'), (2, 'uuid-def')])
    assert cur.executemany.called
    sql = cur.executemany.call_args[0][0]
    assert 'crawler_qqmusic_singer_songs' in sql
  • Step 2: 运行测试确认失败
pytest tests/test_writer.py -v

期望:FAILED(ImportError)

  • Step 3: 实现 etl_to_crawler/writer.py
import uuid
import json
import psycopg2.extras


# ─── QQ Music ────────────────────────────────────────────────────────────────

def upsert_qq_singers(cur, singers: list[dict]) -> None:
    if not singers:
        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())
        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'),
    ) for s in singers]
    cur.executemany(sql, rows)


def upsert_qq_albums(cur, albums: list[dict]) -> None:
    if not albums:
        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())
        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'),
    ) for a in albums]
    cur.executemany(sql, rows)


def upsert_qq_songs(cur, songs: list[dict]) -> None:
    if not songs:
        return
    sql = """
        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())
        ON CONFLICT (platform_song_id) DO NOTHING
    """
    rows = [(
        s['song_uuid'], s['platform_song_id'], s['mid'], s.get('album_id'),
        s.get('cover', ''), s.get('title', ''), s.get('name', ''), s.get('duration', 0) or 0,
        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', '[]'),
    ) for s in songs]
    cur.executemany(sql, rows)


def upsert_qq_singer_songs(cur, pairs: list[tuple]) -> None:
    """pairs: [(singer_id, song_uuid), ...]"""
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_qqmusic_singer_songs (id, singer_id, song_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, song_id) DO NOTHING
    """
    rows = [(str(uuid.uuid4()), singer_id, song_uuid) for singer_id, song_uuid in pairs]
    cur.executemany(sql, rows)


def upsert_qq_singer_albums(cur, pairs: list[tuple]) -> None:
    """pairs: [(singer_id, album_id), ...]"""
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_qqmusic_singer_albums (id, singer_id, album_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, album_id) DO NOTHING
    """
    rows = [(str(uuid.uuid4()), singer_id, album_id) for singer_id, album_id in pairs]
    cur.executemany(sql, rows)


# ─── Kugou ───────────────────────────────────────────────────────────────────

def upsert_kugou_singers(cur, singers: list[dict]) -> None:
    if not singers:
        return
    sql = """
        INSERT INTO crawler_kugou_singers
            (id, name, avatar, sex, area, "index", intro, home_url, created_at, updated_at)
        VALUES (%s, %s, %s, %s::kugou_singer_sex, %s::kugou_singer_area, %s::kugou_singer_index, %s, %s, NOW(), NOW())
        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', ''),
    ) for s in singers]
    cur.executemany(sql, rows)


def upsert_kugou_albums(cur, albums: list[dict]) -> None:
    if not albums:
        return
    sql = """
        INSERT INTO crawler_kugou_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())
        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'),
    ) for a in albums]
    cur.executemany(sql, rows)


def upsert_kugou_songs(cur, songs: list[dict]) -> None:
    if not songs:
        return
    sql = """
        INSERT INTO crawler_kugou_songs
            (id, platform_song_id, hash, album_audio_id, 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, %s::jsonb, 0, NOW(), NOW())
        ON CONFLICT (platform_song_id) DO NOTHING
    """
    rows = [(
        s['song_uuid'], s['platform_song_id'], s.get('hash', ''),
        s.get('album_audio_id', 0) or 0, s.get('album_id'),
        s.get('cover', ''), s.get('title', ''), s.get('name', ''), s.get('duration', 0) or 0,
        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', '[]'),
    ) for s in songs]
    cur.executemany(sql, rows)


def upsert_kugou_singer_songs(cur, pairs: list[tuple]) -> None:
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_kugou_singer_songs (id, singer_id, song_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, song_id) DO NOTHING
    """
    cur.executemany(sql, [(str(uuid.uuid4()), s, sg) for s, sg in pairs])


def upsert_kugou_singer_albums(cur, pairs: list[tuple]) -> None:
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_kugou_singer_albums (id, singer_id, album_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, album_id) DO NOTHING
    """
    cur.executemany(sql, [(str(uuid.uuid4()), s, a) for s, a in pairs])


# ─── Netease ─────────────────────────────────────────────────────────────────

def upsert_netease_singers(cur, singers: list[dict]) -> None:
    if not singers:
        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())
        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'),
    ) for s in singers]
    cur.executemany(sql, rows)


def upsert_netease_albums(cur, albums: list[dict]) -> None:
    if not albums:
        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())
        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'),
    ) for a in albums]
    cur.executemany(sql, rows)


def upsert_netease_songs(cur, songs: list[dict]) -> None:
    if not songs:
        return
    sql = """
        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, singers, status, created_at, updated_at)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW())
        ON CONFLICT (platform_song_id) DO NOTHING
    """
    rows = [(
        s['song_uuid'], s['platform_song_id'], s.get('album_id'),
        s.get('cover', ''), s.get('title', ''), s.get('name', ''), s.get('duration', 0) or 0,
        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', '[]'),
    ) for s in songs]
    cur.executemany(sql, rows)


def upsert_netease_singer_songs(cur, pairs: list[tuple]) -> None:
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_netease_singer_songs (id, singer_id, song_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, song_id) DO NOTHING
    """
    cur.executemany(sql, [(str(uuid.uuid4()), s, sg) for s, sg in pairs])


def upsert_netease_singer_albums(cur, pairs: list[tuple]) -> None:
    if not pairs:
        return
    sql = """
        INSERT INTO crawler_netease_singer_albums (id, singer_id, album_id)
        VALUES (%s, %s, %s)
        ON CONFLICT (singer_id, album_id) DO NOTHING
    """
    cur.executemany(sql, [(str(uuid.uuid4()), s, a) for s, a in pairs])
  • Step 4: 运行测试确认通过
pytest tests/test_writer.py -v

期望:3 个测试全部 PASSED

  • Step 5: Commit
git add etl_to_crawler/writer.py tests/test_writer.py
git commit -m "feat(etl): 实现三平台 PostgreSQL 写入函数"

Task 7: runner.py + run_etl.py — 批次编排与主入口

Files:

  • Create: etl_to_crawler/runner.py
  • Create: run_etl.py

Interfaces:

  • Consumes: reader、spider、writer、oss、lyric、connections 全部模块
  • Produces: runner.run(platforms: list[str]) -> None

  • Step 1: 实现 etl_to_crawler/runner.py

import uuid
import json
import logging
from tqdm import tqdm

from .config import PLATFORM_QQ, PLATFORM_KUGOU, PLATFORM_NETEASE, BATCH_SIZE, OSS_CONFIG
from .connections import get_hk_songs_conn, get_source_conn, get_spider_conn, get_pg_conn, get_oss_bucket
from .reader import iter_hk_songs_batches, fetch_platform_records
from .spider import (
    fetch_qq_songs, fetch_qq_singers,
    fetch_kugou_songs, fetch_kugou_singers,
    fetch_netease_songs, fetch_netease_singers,
)
from .writer import (
    upsert_qq_singers, upsert_qq_albums, upsert_qq_songs,
    upsert_qq_singer_songs, upsert_qq_singer_albums,
    upsert_kugou_singers, upsert_kugou_albums, upsert_kugou_songs,
    upsert_kugou_singer_songs, upsert_kugou_singer_albums,
    upsert_netease_singers, upsert_netease_albums, upsert_netease_songs,
    upsert_netease_singer_songs, upsert_netease_singer_albums,
)
from .oss import transfer_url, build_oss_key
from .lyric import strip_timestamps

logging.basicConfig(level=logging.INFO, format='%(asctime)s %(levelname)s %(message)s')
log = logging.getLogger(__name__)


def _safe_transfer(url, oss_key, bucket, base_url):
    try:
        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,不阻断流程


def _process_qq(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
    mid = pr['platform_unique_key']
    songs_map = fetch_qq_songs(spider_conn, [mid])
    if mid not in songs_map:
        return
    sp = songs_map[mid]
    song_id_int = sp['id']

    singers_map = fetch_qq_singers(spider_conn, [song_id_int])
    singer_list = singers_map.get(song_id_int, [])

    # OSS 转移
    audio_url = _safe_transfer(
        hk_row['audio_url'],
        build_oss_key('qq', 'audio', mid + '.mp3'),
        bucket, base_url
    )
    cover_url = _safe_transfer(
        hk_row.get('cover_url') or sp.get('cover', ''),
        build_oss_key('qq', 'cover', str(song_id_int) + '.jpg'),
        bucket, base_url
    )
    album_cover = ''
    if sp.get('album_id') and sp.get('album_cover'):
        album_cover = _safe_transfer(
            sp['album_cover'],
            build_oss_key('qq', 'album', str(sp['album_id']) + '.jpg'),
            bucket, base_url
        )

    # 歌手头像转移 + 写入 singers
    singer_rows = []
    for sg in singer_list:
        avatar = _safe_transfer(
            sg.get('avatar', ''),
            build_oss_key('qq', 'singer', sg['mid'] + '.jpg'),
            bucket, base_url
        )
        singer_rows.append({**sg, 'avatar': avatar})
    upsert_qq_singers(pg_cur, singer_rows)

    # singers JSONB
    singers_json = json.dumps([{
        'name': sg['name'],
        'singer_id': sg['singer_id'],
        'platform_singer_id': sg['mid'],
    } for sg in singer_list], ensure_ascii=False)

    # 写入 album
    album_id = None
    if sp.get('album_id'):
        album_id = sp['album_id']
        upsert_qq_albums(pg_cur, [{
            'id': sp['album_id'], 'mid': sp.get('album_mid', ''),
            'cover': album_cover, 'title': sp.get('album_title', ''),
            'intro': sp.get('album_intro'), 'type': sp.get('album_type', ''),
            'company_id': sp.get('company_id', 0), 'company': sp.get('company', ''),
            'is_owner': sp.get('is_owner', 0), 'published_at': sp.get('album_published_at'),
        }])

    # 写入 song
    song_uuid = str(uuid.uuid4())
    upsert_qq_songs(pg_cur, [{
        'song_uuid': song_uuid,
        'platform_song_id': int(pr['platform_mid']) if pr.get('platform_mid') else song_id_int,
        'mid': mid,
        'album_id': album_id,
        'cover': cover_url,
        'title': sp.get('title', hk_row['name']),
        'name': hk_row['name'],
        'duration': sp.get('duration', 0) or 0,
        'lyric': strip_timestamps(sp.get('lyric')),
        '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,
        'lyric_url': hk_row.get('lyrics_url'),
        'platform_index_url': sp.get('platform_index_url'),
        'published_at': sp.get('published_at') or hk_row.get('issue_time'),
        'singers_json': singers_json,
    }])

    # singer_songs / singer_albums
    upsert_qq_singer_songs(pg_cur, [(sg['singer_id'], song_uuid) for sg in singer_list])
    if album_id:
        upsert_qq_singer_albums(pg_cur, [(sg['singer_id'], album_id) for sg in singer_list])


def _process_kugou(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
    song_id = int(pr['platform_unique_key'])
    songs_map = fetch_kugou_songs(spider_conn, [song_id])
    if song_id not in songs_map:
        return
    sp = songs_map[song_id]

    singers_map = fetch_kugou_singers(spider_conn, [song_id])
    singer_list = singers_map.get(song_id, [])

    audio_url = _safe_transfer(
        hk_row['audio_url'],
        build_oss_key('kugou', 'audio', str(song_id) + '.mp3'),
        bucket, base_url
    )
    cover_url = _safe_transfer(
        hk_row.get('cover_url') or sp.get('cover', ''),
        build_oss_key('kugou', 'cover', str(song_id) + '.jpg'),
        bucket, base_url
    )
    album_cover = ''
    if sp.get('album_id') and sp.get('album_cover'):
        album_cover = _safe_transfer(
            sp['album_cover'],
            build_oss_key('kugou', 'album', str(sp['album_id']) + '.jpg'),
            bucket, base_url
        )

    singer_rows = []
    for sg in singer_list:
        avatar = _safe_transfer(
            sg.get('avatar', ''),
            build_oss_key('kugou', 'singer', str(sg['singer_id']) + '.jpg'),
            bucket, base_url
        )
        singer_rows.append({**sg, 'id': sg['singer_id'], 'avatar': avatar})
    upsert_kugou_singers(pg_cur, singer_rows)

    singers_json = json.dumps([{
        'name': sg['name'],
        'singer_id': sg['singer_id'],
        'platform_singer_id': str(sg['singer_id']),
    } for sg in singer_list], ensure_ascii=False)

    album_id = None
    if sp.get('album_id'):
        album_id = sp['album_id']
        upsert_kugou_albums(pg_cur, [{
            'id': sp['album_id'], 'cover': album_cover,
            'title': sp.get('album_title', ''), 'intro': sp.get('album_intro'),
            'type': sp.get('album_type', ''), 'company_id': sp.get('company_id', 0),
            'company': sp.get('company', ''), 'is_owner': sp.get('is_owner', 0),
            'published_at': sp.get('album_published_at'),
        }])

    song_uuid = str(uuid.uuid4())
    upsert_kugou_songs(pg_cur, [{
        'song_uuid': song_uuid,
        'platform_song_id': song_id,
        'hash': sp.get('hid', pr.get('platform_mid', '')),
        'album_audio_id': sp.get('album_audio_id') or pr.get('album_audio_id') or 0,
        'album_id': album_id,
        'cover': cover_url,
        'title': sp.get('title', hk_row['name']),
        'name': hk_row['name'],
        'duration': sp.get('duration', 0) or 0,
        'lyric': strip_timestamps(sp.get('lyric')),
        '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,
        'lyric_url': hk_row.get('lyrics_url'),
        'platform_index_url': sp.get('platform_index_url'),
        'published_at': sp.get('published_at') or hk_row.get('issue_time'),
        'singers_json': singers_json,
    }])

    upsert_kugou_singer_songs(pg_cur, [(sg['singer_id'], song_uuid) for sg in singer_list])
    if album_id:
        upsert_kugou_singer_albums(pg_cur, [(sg['singer_id'], album_id) for sg in singer_list])


def _process_netease(hk_row: dict, pr: dict, spider_conn, pg_cur, bucket, base_url):
    song_id = int(pr['platform_unique_key'])
    songs_map = fetch_netease_songs(spider_conn, [song_id])
    if song_id not in songs_map:
        return
    sp = songs_map[song_id]

    singers_map = fetch_netease_singers(spider_conn, [song_id])
    singer_list = singers_map.get(song_id, [])

    audio_url = _safe_transfer(
        hk_row['audio_url'],
        build_oss_key('netease', 'audio', str(song_id) + '.mp3'),
        bucket, base_url
    )
    cover_url = _safe_transfer(
        hk_row.get('cover_url') or sp.get('cover', ''),
        build_oss_key('netease', 'cover', str(song_id) + '.jpg'),
        bucket, base_url
    )
    album_cover = ''
    if sp.get('album_id') and sp.get('album_cover'):
        album_cover = _safe_transfer(
            sp['album_cover'],
            build_oss_key('netease', 'album', str(sp['album_id']) + '.jpg'),
            bucket, base_url
        )

    singer_rows = []
    for sg in singer_list:
        avatar = _safe_transfer(
            sg.get('avatar', ''),
            build_oss_key('netease', 'singer', str(sg['singer_id']) + '.jpg'),
            bucket, base_url
        )
        singer_rows.append({**sg, 'id': sg['singer_id'], 'avatar': avatar})
    upsert_netease_singers(pg_cur, singer_rows)

    singers_json = json.dumps([{
        'name': sg['name'],
        'singer_id': sg['singer_id'],
        'platform_singer_id': str(sg['singer_id']),
    } for sg in singer_list], ensure_ascii=False)

    album_id = None
    if sp.get('album_id'):
        album_id = sp['album_id']
        upsert_netease_albums(pg_cur, [{
            'id': sp['album_id'], 'cover': album_cover,
            'title': sp.get('album_title', ''), 'intro': sp.get('album_intro'),
            'type': sp.get('album_type', ''), 'company_id': sp.get('company_id', 0),
            'company': sp.get('company', ''), 'is_owner': sp.get('is_owner', 0),
            'published_at': sp.get('album_published_at'),
        }])

    song_uuid = str(uuid.uuid4())
    upsert_netease_songs(pg_cur, [{
        'song_uuid': song_uuid,
        'platform_song_id': song_id,
        'album_id': album_id,
        'cover': cover_url,
        'title': sp.get('title', hk_row['name']),
        'name': hk_row['name'],
        'duration': sp.get('duration', 0) or 0,
        'lyric': strip_timestamps(sp.get('lyric')),
        '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,
        'lyric_url': hk_row.get('lyrics_url'),
        'platform_index_url': sp.get('platform_index_url'),
        'published_at': sp.get('published_at') or hk_row.get('issue_time'),
        'singers_json': singers_json,
    }])

    upsert_netease_singer_songs(pg_cur, [(sg['singer_id'], song_uuid) for sg in singer_list])
    if album_id:
        upsert_netease_singer_albums(pg_cur, [(sg['singer_id'], album_id) for sg in singer_list])


_PROCESSORS = {
    PLATFORM_QQ: _process_qq,
    PLATFORM_KUGOU: _process_kugou,
    PLATFORM_NETEASE: _process_netease,
}


def run(platforms: list[str]) -> None:
    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']

    total_ok = total_err = 0

    try:
        for batch in tqdm(iter_hk_songs_batches(hk_conn, BATCH_SIZE), desc='batches'):
            song_ids = [int(r['source_song_id']) for r in batch if r.get('source_song_id')]
            platform_records = fetch_platform_records(src_conn, song_ids)

            # index platform records by source_song_id
            pr_by_song: dict[int, list] = {}
            for pr in platform_records:
                if pr['platform'] in platforms:
                    pr_by_song.setdefault(int(pr['source_song_id']), []).append(pr)

            with pg_conn.cursor() as pg_cur:
                for hk_row in batch:
                    src_id = int(hk_row['source_song_id']) if hk_row.get('source_song_id') else None
                    if not src_id or src_id not in pr_by_song:
                        continue
                    for pr in pr_by_song[src_id]:
                        processor = _PROCESSORS.get(pr['platform'])
                        if not processor:
                            continue
                        try:
                            processor(hk_row, pr, spider_conn, pg_cur, bucket, base_url)
                            total_ok += 1
                        except Exception as e:
                            log.error("Error processing song %s platform %s: %s",
                                      hk_row.get('name'), pr['platform'], e)
                            total_err += 1
                pg_conn.commit()

    finally:
        hk_conn.close()
        src_conn.close()
        spider_conn.close()
        pg_conn.close()

    log.info("Done. OK=%d ERR=%d", total_ok, total_err)
  • Step 2: 创建 run_etl.py
#!/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 run

PLATFORM_MAP = {
    'qq': PLATFORM_QQ,
    'kugou': PLATFORM_KUGOU,
    'netease': PLATFORM_NETEASE,
}

if __name__ == '__main__':
    parser = argparse.ArgumentParser(description='ETL: hk_songs_test → crawler_dev')
    parser.add_argument('--platform', default='all',
                        choices=['qq', 'kugou', 'netease', 'all'],
                        help='要导入的平台(默认 all)')
    args = parser.parse_args()

    if args.platform == 'all':
        platforms = PLATFORMS
    else:
        platforms = [PLATFORM_MAP[args.platform]]

    print(f"Starting ETL for platforms: {platforms}")
    run(platforms)
  • Step 3: 冒烟测试——只跑前 10 条

临时修改 config.py 中 BATCH_SIZE = 10,再运行:

source .venv/bin/activate
python run_etl.py --platform qq

检查:

psql "postgresql://data_test:z1n_z%23JHW0@data-center.rwlb.rds.aliyuncs.com/crawler_dev?sslmode=prefer" \
  -c "SELECT count(*) FROM crawler_qqmusic_songs;" \
  -c "SELECT count(*) FROM crawler_qqmusic_singers;" \
  -c "SELECT count(*) FROM crawler_qqmusic_albums;"

期望:三张表均有 > 0 的行数,且无报错。

恢复 BATCH_SIZE = 500

  • Step 4: 全量运行三平台
source .venv/bin/activate
python run_etl.py --platform all 2>&1 | tee etl_run.log

完成后检查:

psql "postgresql://data_test:z1n_z%23JHW0@data-center.rwlb.rds.aliyuncs.com/crawler_dev?sslmode=prefer" -c "
SELECT 'qqmusic_songs' AS tbl, count(*) FROM crawler_qqmusic_songs
UNION ALL SELECT 'qqmusic_singers', count(*) FROM crawler_qqmusic_singers
UNION ALL SELECT 'qqmusic_albums', count(*) FROM crawler_qqmusic_albums
UNION ALL SELECT 'kugou_songs', count(*) FROM crawler_kugou_songs
UNION ALL SELECT 'kugou_singers', count(*) FROM crawler_kugou_singers
UNION ALL SELECT 'kugou_albums', count(*) FROM crawler_kugou_albums
UNION ALL SELECT 'netease_songs', count(*) FROM crawler_netease_songs
UNION ALL SELECT 'netease_singers', count(*) FROM crawler_netease_singers
UNION ALL SELECT 'netease_albums', count(*) FROM crawler_netease_albums;
"

期望:9 张主表均有数据。

  • Step 5: Commit
git add etl_to_crawler/runner.py run_etl.py
git commit -m "feat(etl): 实现批次编排与 CLI 入口,完成全量导入"