db_reader.py 2.8 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394
  1. """
  2. 数据库读取模块(只读)
  3. - PostgreSQL: 读取 cards_master_v2 卡牌元数据
  4. - ClickHouse: 读取 card_transactions_unified_sync 交易数据
  5. """
  6. import sys
  7. sys.stdout.reconfigure(encoding="utf-8")
  8. def read_card_master(pg_config, table, filters=None, limit=None):
  9. """
  10. 读取 cards_master_v2。
  11. Args:
  12. filters: dict, 例 {"language":["简中"], "year":["1996","1999"]}
  13. limit: 限制返回条数(调试用)
  14. Returns:
  15. list of dict: [{card_id, card_name_ch, language, year, card_no, pg_label, img_url}, ...]
  16. """
  17. import psycopg2
  18. where = ["img_url IS NOT NULL", "img_url != ''"]
  19. params = []
  20. if filters:
  21. for key, vals in filters.items():
  22. if vals:
  23. placeholders = ",".join(["%s"] * len(vals))
  24. where.append(f"{key} IN ({placeholders})")
  25. params.extend(vals)
  26. sql = f"SELECT card_id, card_name_ch, language, year, card_no, pg_label, img_url " \
  27. f"FROM {table} WHERE {' AND '.join(where)} ORDER BY card_id"
  28. if limit:
  29. sql += f" LIMIT {int(limit)}"
  30. conn = psycopg2.connect(**pg_config)
  31. try:
  32. cur = conn.cursor()
  33. cur.execute(sql, params)
  34. cols = [d[0] for d in cur.description]
  35. rows = [dict(zip(cols, r)) for r in cur.fetchall()]
  36. cur.close()
  37. return rows
  38. finally:
  39. conn.close()
  40. def read_transactions(ch_config, table, limit=None, where=None):
  41. """
  42. 读取交易数据。
  43. Args:
  44. limit: 限制条数
  45. where: 额外 WHERE 条件字符串
  46. Returns:
  47. list of dict: [{tx_id, card_id, platform, title, image_uri, cos_image_url, ...}, ...]
  48. """
  49. import clickhouse_connect
  50. sql = f"SELECT * FROM {table}"
  51. if where:
  52. sql += f" WHERE {where}"
  53. sql += " ORDER BY sold_date DESC"
  54. if limit:
  55. sql += f" LIMIT {int(limit)}"
  56. client = clickhouse_connect.get_client(**ch_config)
  57. try:
  58. result = client.query(sql)
  59. cols = result.column_names
  60. return [dict(zip(cols, r)) for r in result.result_rows]
  61. finally:
  62. client.close()
  63. def count_card_master(pg_config, table, filters=None):
  64. """统计满足条件的卡牌数量"""
  65. import psycopg2
  66. where = ["img_url IS NOT NULL", "img_url != ''"]
  67. params = []
  68. if filters:
  69. for key, vals in filters.items():
  70. if vals:
  71. placeholders = ",".join(["%s"] * len(vals))
  72. where.append(f"{key} IN ({placeholders})")
  73. params.extend(vals)
  74. sql = f"SELECT COUNT(*) FROM {table} WHERE {' AND '.join(where)}"
  75. conn = psycopg2.connect(**pg_config)
  76. try:
  77. cur = conn.cursor()
  78. cur.execute(sql, params)
  79. cnt = cur.fetchone()[0]
  80. cur.close()
  81. return cnt
  82. finally:
  83. conn.close()