From a494f4e72985edbc277a2e23f199a6a0d4c6e8ca Mon Sep 17 00:00:00 2001 From: shellway-pc <413209390@qq.com> Date: Wed, 5 Aug 2026 22:21:46 +0800 Subject: [PATCH] =?UTF-8?q?fix=EF=BC=9A=E4=BF=AE=E5=A4=8D=E4=BA=86?= =?UTF-8?q?=E8=B4=A2=E5=8A=A1=E6=95=B0=E6=8D=AE=E5=AF=BC=E5=85=A5=E8=BF=87?= =?UTF-8?q?=E7=A8=8B=E4=B8=AD=E5=87=BA=E7=8E=B0=E9=87=8D=E5=A4=8D=E5=88=97?= =?UTF-8?q?=E7=9A=84bug=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- quantitative_data/importer.py | 34 +++++++++++++++++++++++++++++++--- 1 file changed, 31 insertions(+), 3 deletions(-) diff --git a/quantitative_data/importer.py b/quantitative_data/importer.py index 650c2a5..ab83f76 100644 --- a/quantitative_data/importer.py +++ b/quantitative_data/importer.py @@ -150,9 +150,6 @@ def batch_insert(table_name: str, df: pd.DataFrame, conn, conflict_columns: List logger.warning(f" {table_name}: 过滤后无可用列,跳过") return 0 - # pd.NaT / pd.NaT / numpy NaN 等无法被 psycopg2 识别,统一替换为 Python None - df = df.where(pd.notna(df), None) - columns = list(df.columns) # 过滤后验证冲突列仍存在 @@ -163,6 +160,37 @@ def batch_insert(table_name: str, df: pd.DataFrame, conn, conflict_columns: List ) return 0 + # ---- 按冲突列去重 ---- + # PostgreSQL 的 ON CONFLICT DO UPDATE 不允许同一命令中出现重复冲突键: + # "ON CONFLICT DO UPDATE command cannot affect row a second time" + # Tushare 财务接口 (income/balancesheet/cashflow/fina_indicator) 对同一股票 + # 同一报告期可能返回多行数据(如多次公告修正,ann_date 不同但冲突键相同), + # 必须先在批内去重。 + # + # 为保证保留的是"最新公告"的数据而非仅依赖 Tushare 返回顺序: + # 若存在公告日期列 (ann_date / f_ann_date),先按公告日期升序排序, + # 再 drop_duplicates(keep="last") 即可稳定保留最新一条 (公告日期最大), + # 且缺失公告日期的行 (NaT) 会排在最后,仅当无公告日期时才被保留。 + # + # 注意:去重必须在此处 (NaN->None 替换之前) 执行, + # 此时日期列仍为 datetime64 类型,sort_values(na_position="last") + # 能正确处理 NaT;若在替换之后排序,object 类型混合日期/None 排序不可靠。 + dedup_cols = [c for c in conflict_columns if c in columns] + before_dedup = len(df) + date_cols = [c for c in ["f_ann_date", "ann_date"] if c in columns] + if date_cols: + df = df.sort_values(date_cols, na_position="last") + df = df.drop_duplicates(subset=dedup_cols, keep="last") + after_dedup = len(df) + if after_dedup < before_dedup: + logger.warning( + f" {table_name}: 检测到 {before_dedup - after_dedup} 行重复冲突键," + f"已保留最新公告记录 (去重后 {after_dedup} 行)" + ) + + # pd.NaT / numpy NaN 等无法被 psycopg2 识别,统一替换为 Python None + df = df.where(pd.notna(df), None) + rows = [tuple(row) for row in df.itertuples(index=False)] # 构建 ON CONFLICT 子句(使用 sql.Identifier 防止注入/语法错误)