""" 数据库读取模块(只读) - PostgreSQL: 读取 cards_master_v2 卡牌元数据 - ClickHouse: 读取 card_transactions_unified_sync 交易数据 """ import sys sys.stdout.reconfigure(encoding="utf-8") def read_card_master(pg_config, table, filters=None, limit=None): """ 读取 cards_master_v2。 Args: filters: dict, 例 {"language":["简中"], "year":["1996","1999"]} limit: 限制返回条数(调试用) Returns: list of dict: [{card_id, card_name_ch, language, year, card_no, pg_label, img_url}, ...] """ import psycopg2 where = ["img_url IS NOT NULL", "img_url != ''"] params = [] if filters: for key, vals in filters.items(): if vals: placeholders = ",".join(["%s"] * len(vals)) where.append(f"{key} IN ({placeholders})") params.extend(vals) sql = f"SELECT card_id, card_name_ch, language, year, card_no, pg_label, img_url " \ f"FROM {table} WHERE {' AND '.join(where)} ORDER BY card_id" if limit: sql += f" LIMIT {int(limit)}" conn = psycopg2.connect(**pg_config) try: cur = conn.cursor() cur.execute(sql, params) cols = [d[0] for d in cur.description] rows = [dict(zip(cols, r)) for r in cur.fetchall()] cur.close() return rows finally: conn.close() def read_transactions(ch_config, table, limit=None, where=None): """ 读取交易数据。 Args: limit: 限制条数 where: 额外 WHERE 条件字符串 Returns: list of dict: [{tx_id, card_id, platform, title, image_uri, cos_image_url, ...}, ...] """ import clickhouse_connect sql = f"SELECT * FROM {table}" if where: sql += f" WHERE {where}" sql += " ORDER BY sold_date DESC" if limit: sql += f" LIMIT {int(limit)}" client = clickhouse_connect.get_client(**ch_config) try: result = client.query(sql) cols = result.column_names return [dict(zip(cols, r)) for r in result.result_rows] finally: client.close() def count_card_master(pg_config, table, filters=None): """统计满足条件的卡牌数量""" import psycopg2 where = ["img_url IS NOT NULL", "img_url != ''"] params = [] if filters: for key, vals in filters.items(): if vals: placeholders = ",".join(["%s"] * len(vals)) where.append(f"{key} IN ({placeholders})") params.extend(vals) sql = f"SELECT COUNT(*) FROM {table} WHERE {' AND '.join(where)}" conn = psycopg2.connect(**pg_config) try: cur = conn.cursor() cur.execute(sql, params) cnt = cur.fetchone()[0] cur.close() return cnt finally: conn.close()