writer.py 14.3 KB
import uuid


def insert_yinyan_song_records(cur, records: list[dict]) -> None:
    """Initialize yinyan_song_records rows before crawler import."""
    if not records:
        return
    cur.executemany(
        """
        INSERT INTO yinyan_song_records (song_id, record_id, platform, is_yinyan_push)
        VALUES (%s, %s, %s, FALSE)
        ON CONFLICT (song_id, record_id) DO UPDATE
            SET platform = EXCLUDED.platform
            WHERE yinyan_song_records.platform IS NULL
        """,
        [(r['song_id'], r['record_id'], r['platform']) for r in records],
    )


def fetch_existing_yinyan_song_ids(cur) -> set[int]:
    """返回已完成初始化 platform 的 song_id 集合。"""
    cur.execute("SELECT DISTINCT song_id FROM yinyan_song_records WHERE platform IS NOT NULL")
    return {row[0] for row in cur.fetchall()}


def fetch_pending_yinyan_song_records(cur, limit: int) -> list[dict]:
    cur.execute(
        """
        SELECT song_id, record_id, platform
        FROM yinyan_song_records
        WHERE is_yinyan_push = FALSE
          AND platform IS NOT NULL
        ORDER BY song_id
        LIMIT %s
        """,
        (limit,),
    )
    rows = cur.fetchall()
    return [{'song_id': row[0], 'record_id': row[1], 'platform': row[2]} for row in rows]


def fetch_yinyan_records_missing_platform(cur, limit: int) -> list[dict]:
    cur.execute(
        """
        SELECT song_id, record_id
        FROM yinyan_song_records
        WHERE platform IS NULL
        ORDER BY song_id
        LIMIT %s
        """,
        (limit,),
    )
    rows = cur.fetchall()
    return [{'song_id': row[0], 'record_id': row[1]} for row in rows]


def update_yinyan_record_platforms(cur, records: list[dict]) -> None:
    if not records:
        return
    values_sql = ', '.join(['(%s::bigint, %s::bigint, %s::varchar)'] * len(records))
    params = []
    for r in records:
        params.extend([r['song_id'], r['record_id'], r['platform']])
    cur.execute(
        f"""
        UPDATE yinyan_song_records AS ysr
        SET platform = v.platform
        FROM (VALUES {values_sql}) AS v(song_id, record_id, platform)
        WHERE ysr.song_id = v.song_id
          AND ysr.record_id = v.record_id
          AND ysr.platform IS NULL
        """,
        tuple(params),
    )


def upsert_yinyan_song_records(cur, records: list[dict]) -> None:
    """Mark pre-initialized yinyan song-record rows as pushed to crawler.
    Matches only on song_id so singer-fallback can use a different record_id than initialized.
    """
    if not records:
        return
    values_sql = ', '.join(['(%s::bigint, %s::bigint, %s::varchar, %s::bigint)'] * len(records))
    params = []
    for r in records:
        params.extend([r['song_id'], r['record_id'], r['platform'], r['platform_song_id']])
    cur.execute(
        f"""
        UPDATE yinyan_song_records AS ysr
        SET platform = v.platform,
            platform_song_id = v.platform_song_id,
            record_id = v.record_id,
            is_yinyan_push = TRUE
        FROM (VALUES {values_sql}) AS v(song_id, record_id, platform, platform_song_id)
        WHERE ysr.song_id = v.song_id
          AND ysr.is_yinyan_push = FALSE
        """,
        tuple(params),
    )


# ─── 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s::singer_sex, %s::singer_area, %s::singer_index, %s, %s,
                NOW(), NOW(), %s, %s::json)
        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'),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
        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'),
        a.get('provider_name'), a.get('crawler_source_data'),
    ) 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, album, singers, status, created_at, updated_at,
             version, audio_md5, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW(), %s, %s, %s, %s::json)
        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('album_json'), s.get('singers_json', '[]'),
        s.get('version'),
        s.get('audio_md5'),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s::kugou_singer_sex, %s::kugou_singer_area, %s::kugou_singer_index, %s, %s,
                NOW(), NOW(), %s, %s::json)
        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', ''),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
        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'),
        a.get('provider_name'), a.get('crawler_source_data'),
    ) 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,
             version, audio_md5, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, 0, NOW(), NOW(), %s, %s, %s, %s::json)
        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', '[]'),
        s.get('version'),
        s.get('audio_md5'),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s::netease_singer_sex, %s::netease_singer_area, %s::netease_singer_index, %s, %s,
                NOW(), NOW(), %s, %s::json)
        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'),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, NOW(), NOW(), %s, %s::json)
        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'),
        a.get('provider_name'), a.get('crawler_source_data'),
    ) 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, album, singers, status, created_at, updated_at,
             version, audio_md5, provider_name, crawler_source_data)
        VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s::json, %s::jsonb, 0, NOW(), NOW(), %s, %s, %s, %s::json)
        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('album_json'), s.get('singers_json', '[]'),
        s.get('version'),
        s.get('audio_md5'),
        s.get('provider_name'), s.get('crawler_source_data'),
    ) 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])