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 防止注入/语法错误)