Commit fac42d3e fac42d3e63240e369ede39d69b7e6db7c38c889e by 沈秋雨

docs: 补充相关文档

1 parent bde8dd5f
# 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_HOST``CRAWLER_DB_PORT``CRAWLER_DB_USER``CRAWLER_DB_PASSWORD``CRAWLER_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**
```bash
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`**
```python
```
(空文件)
- [ ] **Step 4: 创建 `etl_to_crawler/config.py`**
```python
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`**
```python
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: 验证所有连接**
```bash
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**
```bash
pip freeze > requirements.txt
```
- [ ] **Step 8: Commit**
```bash
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: 写测试**
```python
# 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: 运行测试确认失败**
```bash
source .venv/bin/activate && pytest tests/test_lyric.py -v
```
期望:FAILED(ImportError 或 NameError)
- [ ] **Step 3: 实现 `etl_to_crawler/lyric.py`**
```python
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: 运行测试确认通过**
```bash
pytest tests/test_lyric.py -v
```
期望:5 个测试全部 PASSED
- [ ] **Step 5: Commit**
```bash
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: 写测试**
```python
# 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: 运行测试确认失败**
```bash
pytest tests/test_oss.py -v
```
期望:FAILED(ImportError)
- [ ] **Step 3: 实现 `etl_to_crawler/oss.py`**
```python
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: 运行测试确认通过**
```bash
pytest tests/test_oss.py -v
```
期望:4 个测试全部 PASSED
- [ ] **Step 5: Commit**
```bash
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 格式:
```python
{
'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`**
```python
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: 手动验证**
```bash
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**
```bash
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: 写测试**
```python
# 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: 运行测试确认失败**
```bash
pytest tests/test_spider.py -v
```
期望:FAILED(ImportError)
- [ ] **Step 3: 实现 `etl_to_crawler/spider.py`**
```python
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: 运行测试确认通过**
```bash
pytest tests/test_spider.py -v
```
期望:3 个测试全部 PASSED
- [ ] **Step 5: Commit**
```bash
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: 写测试**
```python
# 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: 运行测试确认失败**
```bash
pytest tests/test_writer.py -v
```
期望:FAILED(ImportError)
- [ ] **Step 3: 实现 `etl_to_crawler/writer.py`**
```python
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: 运行测试确认通过**
```bash
pytest tests/test_writer.py -v
```
期望:3 个测试全部 PASSED
- [ ] **Step 5: Commit**
```bash
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`**
```python
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`**
```python
#!/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`,再运行:
```bash
source .venv/bin/activate
python run_etl.py --platform qq
```
检查:
```bash
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: 全量运行三平台**
```bash
source .venv/bin/activate
python run_etl.py --platform all 2>&1 | tee etl_run.log
```
完成后检查:
```bash
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**
```bash
git add etl_to_crawler/runner.py run_etl.py
git commit -m "feat(etl): 实现批次编排与 CLI 入口,完成全量导入"
```
# 数据流文档:hk_songs_test → crawler_dev(三平台)
---
## 一、数据库说明
| 数据库 | 连接 | 用途 |
|---|---|---|
| `new_music_library` | TencentDB (port 60209) | 包含 `hk_songs_test`(源数据) |
| `hikoon-data` | AliCloud RDS (port 3307) | 包含平台关联表 `hk_song_platform``hk_song_and_record``hk_music_record` 等 |
| `hikoon-data-spider` | 同上 AliCloud RDS (port 3307) | 包含各平台爬虫原始表 `media_tencent_songs` 等 |
| `crawler_dev` | AliCloud RDS PostgreSQL | 目标写入库 |
---
## 二、关联链路
```
new_music_library.hk_songs_test.source_song_id
= hikoon-data.hk_song_platform.id
↓ (hk_song_and_record.song_id = hk_song_platform.id)
hikoon-data.hk_song_and_record
↓ (record_id = hk_music_record.id)
hikoon-data.hk_music_record(一首歌对应多条,每个平台/版本各一条)
├── platform (1=QQ, 2=Kugou, 4=Netease)
├── platform_unique_key
│ ├─ QQ → media_tencent_songs.mid(字符串 mid)
│ ├─ Kugou → media_ku_gou_songs.id(数字 ID)
│ └─ Netease → media_netease_songs.id(数字 ID)
└── platform_mid
├─ QQ → media_tencent_songs.id(数字 ID)
└─ Kugou → media_ku_gou_songs.hid(hash 字符串)
```
> **重要**:`hk_songs_test.source_table_name = 'hk_song_platform'`,`source_song_id` 即为 `hk_song_platform.id`,是从 hk_songs_test 回溯到 hikoon-data 的唯一关联键。
>
> **一对多**:一条 hk_songs_test 记录可能在 QQ、酷狗、网易都有对应的 `hk_music_record`,每个平台分别写入各自的 `crawler_xx_songs` 表。
---
## 三、QQ音乐(platform = 1)
**Spider 核心关联**`hk_music_record.platform_unique_key` = `media_tencent_songs.mid`
### crawler_dev.crawler_qqmusic_songs
| 目标字段 | 来源库 | 来源表.字段 | 备注 |
|---|---|---|---|
| `id` | 生成 | — | UUID |
| `platform_song_id` | hikoon-data | `hk_music_record.platform_mid` | QQ 数字歌曲 ID,等同于 `media_tencent_songs.id` |
| `mid` | hikoon-data | `hk_music_record.platform_unique_key` | QQ mid 字符串,等同于 `media_tencent_songs.mid` |
| `title` / `name` | new_music_library | `hk_songs_test.name` | |
| `cover` | new_music_library → OSS | `hk_songs_test.cover_url` | 转换至 archive-dev OSS |
| `duration` | hikoon-data-spider | `media_tencent_songs.duration` | hk_songs_test 无此字段 |
| `lyric` | hikoon-data-spider | `media_tencent_songs.lyric` | 去掉 `[mm:ss.xx]` 时间戳后存储 |
| `composer_name` | hikoon-data-spider | `media_tencent_songs.composer_name` | 备选 `hk_songs_test.composer` |
| `lyricist_name` | hikoon-data-spider | `media_tencent_songs.lyricist_name` | 备选 `hk_songs_test.lyricist` |
| `url` | new_music_library → OSS | `hk_songs_test.audio_url` | 统一转换至 archive-dev OSS |
| `lyric_url` | new_music_library | `hk_songs_test.lyrics_url` | 已是 archive-dev OSS .txt,直接使用 |
| `published_at` | hikoon-data-spider | `media_tencent_songs.published_at` | 备选 `hk_songs_test.issue_time` |
| `platform_index_url` | hikoon-data-spider | `media_tencent_songs.platform_index_url` | |
| `album_id` | 关联后写入 | 见下方 albums 说明 | FK → `crawler_qqmusic_albums.id` |
| `singers` | 关联后组装 | 见下方 singers 说明 | JSONB 格式 |
| `album` | hikoon-data-spider | `media_tencent_albums.*` 组成 JSON | 冗余字段 |
### crawler_dev.crawler_qqmusic_singers
**关联路径**`media_tencent_songs.id`(by mid)→ `media_tencent_singer_has_songs.song_id``media_tencent_singers`(ON singer_id)
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_tencent_singers.id` |
| `mid` | `media_tencent_singers.mid` |
| `name` | `media_tencent_singers.name` |
| `avatar` | `media_tencent_singers.avatar` → OSS |
| `sex` | `media_tencent_singers.sex` |
| `area` | `media_tencent_singers.area` |
| `index` | `media_tencent_singers.index` |
| `intro` | `media_tencent_singers.intro` |
| `home_url` | `media_tencent_singers.home_url` |
**`crawler_qqmusic_songs.singers` JSONB 格式**
```json
[{"name": "歌手名", "singer_id": <media_tencent_singers.id>, "platform_singer_id": "<media_tencent_singers.mid>"}]
```
### crawler_dev.crawler_qqmusic_albums
**关联路径**`media_tencent_songs.album_id``media_tencent_albums.id`
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_tencent_albums.id` |
| `mid` | `media_tencent_albums.mid` |
| `cover` | `media_tencent_albums.cover` → OSS |
| `title` | `media_tencent_albums.title` |
| `intro` | `media_tencent_albums.intro` |
| `type` | `media_tencent_albums.type` |
| `company_id` | `media_tencent_albums.company_id` |
| `company` | `media_tencent_albums.company` |
| `is_owner` | `media_tencent_albums.is_owner` |
| `published_at` | `media_tencent_albums.published_at` |
---
## 四、酷狗(platform = 2)
**Spider 核心关联**`hk_music_record.platform_unique_key` = `media_ku_gou_songs.id`
### crawler_dev.crawler_kugou_songs
| 目标字段 | 来源库 | 来源表.字段 | 备注 |
|---|---|---|---|
| `id` | 生成 | — | UUID |
| `platform_song_id` | hikoon-data | `hk_music_record.platform_unique_key` | 酷狗数字歌曲 ID,等同于 `media_ku_gou_songs.id` |
| `hash` | hikoon-data | `hk_music_record.platform_mid` | 酷狗 hash 字符串,等同于 `media_ku_gou_songs.hid` |
| `album_audio_id` | hikoon-data | `hk_music_record.album_audio_id` | |
| `title` / `name` | new_music_library | `hk_songs_test.name` | |
| `cover` | new_music_library → OSS | `hk_songs_test.cover_url` | 转换至 archive-dev OSS |
| `duration` | hikoon-data-spider | `media_ku_gou_songs.duration` | |
| `lyric` | hikoon-data-spider | `media_ku_gou_songs.lyric` | 去时间戳后存储 |
| `composer_name` | hikoon-data-spider | `media_ku_gou_songs.composer_name` | |
| `lyricist_name` | hikoon-data-spider | `media_ku_gou_songs.lyricist_name` | |
| `url` | new_music_library → OSS | `hk_songs_test.audio_url` | 统一转换至 archive-dev OSS |
| `lyric_url` | new_music_library | `hk_songs_test.lyrics_url` | 已是 archive-dev OSS .txt,直接使用 |
| `published_at` | hikoon-data-spider | `media_ku_gou_songs.published_at` | |
| `platform_index_url` | hikoon-data-spider | `media_ku_gou_songs.platform_index_url` | |
| `album_id` | 关联后写入 | 见 albums | FK → `crawler_kugou_albums.id` |
| `singers` | 关联后组装 | 见 singers | JSONB |
### crawler_dev.crawler_kugou_singers
**关联路径**`media_ku_gou_songs.id``media_ku_gou_singer_has_songs.song_id``media_ku_gou_singers`(ON singer_id)
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_ku_gou_singers.id` |
| `name` | `media_ku_gou_singers.name` |
| `avatar` | `media_ku_gou_singers.avatar` → OSS |
| `sex` | `media_ku_gou_singers.sex` |
| `area` | `media_ku_gou_singers.area` |
| `index` | `media_ku_gou_singers.index` |
| `intro` | `media_ku_gou_singers.intro` |
| `home_url` | `media_ku_gou_singers.home_url` |
**`crawler_kugou_songs.singers` JSONB 格式**
```json
[{"name": "歌手名", "singer_id": <media_ku_gou_singers.id>, "platform_singer_id": "<id字符串>"}]
```
### crawler_dev.crawler_kugou_albums
**关联路径**`media_ku_gou_songs.album_id``media_ku_gou_albums.id`
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_ku_gou_albums.id` |
| `cover` | `media_ku_gou_albums.cover` → OSS |
| `title` | `media_ku_gou_albums.title` |
| `intro` | `media_ku_gou_albums.intro` |
| `type` | `media_ku_gou_albums.type` |
| `company_id` | `media_ku_gou_albums.company_id` |
| `company` | `media_ku_gou_albums.company` |
| `is_owner` | `media_ku_gou_albums.is_owner` |
| `published_at` | `media_ku_gou_albums.published_at` |
---
## 五、网易云(platform = 4)
**Spider 核心关联**`hk_music_record.platform_unique_key` = `media_netease_songs.id`
### crawler_dev.crawler_netease_songs
| 目标字段 | 来源库 | 来源表.字段 | 备注 |
|---|---|---|---|
| `id` | 生成 | — | UUID |
| `platform_song_id` | hikoon-data | `hk_music_record.platform_unique_key` | 网易云数字歌曲 ID,等同于 `media_netease_songs.id` |
| `title` / `name` | new_music_library | `hk_songs_test.name` | |
| `cover` | new_music_library → OSS | `hk_songs_test.cover_url` | 转换至 archive-dev OSS |
| `duration` | hikoon-data-spider | `media_netease_songs.duration` | |
| `lyric` | hikoon-data-spider | `media_netease_songs.lyric` | 去时间戳后存储 |
| `composer_name` | hikoon-data-spider | `media_netease_songs.composer_name` | |
| `lyricist_name` | hikoon-data-spider | `media_netease_songs.lyricist_name` | |
| `url` | new_music_library → OSS | `hk_songs_test.audio_url` | 统一转换至 archive-dev OSS |
| `lyric_url` | new_music_library | `hk_songs_test.lyrics_url` | 已是 archive-dev OSS .txt,直接使用 |
| `published_at` | hikoon-data-spider | `media_netease_songs.published_at` | |
| `platform_index_url` | hikoon-data-spider | `media_netease_songs.platform_index_url` | |
| `album_id` | 关联后写入 | 见 albums | FK → `crawler_netease_albums.id` |
| `singers` | 关联后组装 | 见 singers | JSONB |
### crawler_dev.crawler_netease_singers
**关联路径**`media_netease_songs.id``media_netease_singer_has_songs.song_id``media_netease_singers`(ON singer_id)
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_netease_singers.id` |
| `name` | `media_netease_singers.name` |
| `avatar` | `media_netease_singers.avatar` → OSS |
| `sex` | `media_netease_singers.sex` |
| `area` | `media_netease_singers.area` |
| `index` | `media_netease_singers.index` |
| `intro` | `media_netease_singers.intro` |
| `home_url` | `media_netease_singers.home_url` |
**`crawler_netease_songs.singers` JSONB 格式**
```json
[{"name": "歌手名", "singer_id": <media_netease_singers.id>, "platform_singer_id": "<id字符串>"}]
```
### crawler_dev.crawler_netease_albums
**关联路径**`media_netease_songs.album_id``media_netease_albums.id`
| 目标字段 | 来源表.字段 |
|---|---|
| `id` | `media_netease_albums.id` |
| `cover` | `media_netease_albums.cover` → OSS |
| `title` | `media_netease_albums.title` |
| `intro` | `media_netease_albums.intro` |
| `type` | `media_netease_albums.type` |
| `company_id` | `media_netease_albums.company_id` |
| `company` | `media_netease_albums.company` |
| `is_owner` | `media_netease_albums.is_owner` |
| `published_at` | `media_netease_albums.published_at` |
---
## 六、入库顺序与过滤条件
### 过滤条件
| 条件 | 实际字段 | 处理 |
|---|---|---|
| 平台记录存在 | `hk_music_record.platform_unique_key` 不为空 | 关联后若无对应平台记录则跳过 |
| 标题不为空 | `hk_songs_test.name` | 空则跳过 |
| 音频URL不为空 | `hk_songs_test.audio_url` | 空则跳过;有则统一转换至 archive-dev OSS |
| 歌手不为空 | `hk_songs_test.singer` | 文本字段,空则跳过;实际 singers 数据从 spider 关联表获取 |
| 歌词 | `hk_songs_test.lyrics_url` | 已是 archive-dev OSS .txt,直接写入 `lyric_url`;spider 的 `lyric` 字段去时间戳后写入 `lyric` 列 |
### 写入顺序(每首歌,每个平台)
1. **singers** — 按 `id``ON CONFLICT DO NOTHING`
2. **albums** — 按 `id``ON CONFLICT DO NOTHING`
3. **songs** — 填入 `album_id`(FK)和 `singers`(JSONB),按 `platform_song_id` 做唯一约束
4. **singer_songs** — songs 写入后填充,按 `(singer_id, song_id)``ON CONFLICT DO NOTHING`
5. **singer_albums** — albums 写入后填充,按 `(singer_id, album_id)``ON CONFLICT DO NOTHING`
---
## 七、歌手关联表
### crawler_xx_singer_songs
| 目标字段 | 类型 | 来源 |
|---|---|---|
| `id` | uuid | 生成 |
| `singer_id` | int/bigint | `crawler_xx_singers.id`(已写入的歌手主键) |
| `song_id` | uuid | `crawler_xx_songs.id`(已写入的歌曲主键) |
**数据来源**`media_xx_singer_has_songs`(singer_id, song_id),spider 表的 singer_id / song_id 与 PG 各表主键直接对应。
| PG 目标表 | Spider 来源表 | 唯一约束 |
|---|---|---|
| `crawler_qqmusic_singer_songs` | `media_tencent_singer_has_songs` | `(singer_id, song_id)` |
| `crawler_kugou_singer_songs` | `media_ku_gou_singer_has_songs` | `(singer_id, song_id)` |
| `crawler_netease_singer_songs` | `media_netease_singer_has_songs` | `(singer_id, song_id)` |
### crawler_xx_singer_albums
| 目标字段 | 类型 | 来源 |
|---|---|---|
| `id` | uuid | 生成 |
| `singer_id` | int/bigint | `crawler_xx_singers.id`(已写入的歌手主键) |
| `album_id` | bigint | `crawler_xx_albums.id`(已写入的专辑主键) |
**数据来源**`media_xx_singer_has_albums`(singer_id, album_id)。
| PG 目标表 | Spider 来源表 | 唯一约束 |
|---|---|---|
| `crawler_qqmusic_singer_albums` | `media_tencent_singer_has_albums` | `(singer_id, album_id)` |
| `crawler_kugou_singer_albums` | `media_ku_gou_singer_has_albums` | `(singer_id, album_id)` |
| `crawler_netease_singer_albums` | `media_netease_singer_has_albums` | `(singer_id, album_id)` |