使用pandas to_sql()导入CSV至数据库表时忽略错误行的方法
解决pandas to_sql遇错误行导致全量插入失败的方案
方案1:分批次插入(平衡效率与容错)
单行插入太慢、全量插入容错性差,分批次是最优折衷方案。先按较大批次插入,失败后递归拆分批次定位错误行,跳过错误行继续处理:
import pandas as pd from sqlalchemy import create_engine def batch_insert(df, tablename, engine, batch_size=1000): total_rows = len(df) for start in range(0, total_rows, batch_size): end = min(start + batch_size, total_rows) batch = df.iloc[start:end] try: with engine.begin() as conn: batch.to_sql(name=tablename, con=conn, index=False, if_exists='append') print(f"成功插入批次:{start}-{end-1}") except Exception as ex: print(f"批次 {start}-{end-1} 插入失败,尝试拆分: {str(ex)}") # 拆分批次到更小单元,比如每10行一批 if batch_size > 10: batch_insert(batch, tablename, engine, batch_size=10) else: # 小批次仍失败,逐行排查 for idx, row in batch.iterrows(): try: with engine.begin() as conn: pd.DataFrame([row]).to_sql(name=tablename, con=conn, index=False, if_exists='append') print(f"成功插入行 {idx}") except Exception as row_ex: print(f"行 {idx} 插入失败,跳过: {str(row_ex)}") # 使用示例 engine = create_engine("your_database_connection_string") data = pd.read_csv(filename, delimiter='|', skipinitialspace=True, quoting=csv.QUOTE_NONE, quotechar='"', escapechar='\\') df = pd.DataFrame(data) batch_insert(df, tablename, engine)
这个方法的优势是大部分正常数据用大批次快速插入,只有出错的批次才会拆分到小批量甚至单行,整体效率比纯单行插入高很多。
方案2:提前数据校验(从源头避免错误)
根据数据库表的约束(字段类型、长度、格式等),提前在DataFrame中校验并过滤不符合要求的行,把错误行单独保存以便后续排查:
import pandas as pd import csv # 示例:根据数据库表约束定义校验规则 def validate_row(row): # 校验id是否为合法整数 try: int(row['id']) except ValueError: return False, "id不是合法整数" # 校验name字段长度不超过50 if len(str(row['name'])) > 50: return False, "name长度超过限制" # 校验日期格式合法性 try: pd.to_datetime(row['create_date']) except ValueError: return False, "create_date格式错误" return True, "" # 读取数据并执行校验 data = pd.read_csv(filename, delimiter='|', skipinitialspace=True, quoting=csv.QUOTE_NONE, quotechar='"', escapechar='\\') df = pd.DataFrame(data) valid_rows = [] error_rows = [] for idx, row in df.iterrows(): is_valid, msg = validate_row(row) if is_valid: valid_rows.append(row) else: error_rows.append((idx, msg, row)) # 插入有效行 if valid_rows: valid_df = pd.DataFrame(valid_rows) with engine.begin() as conn: valid_df.to_sql(name=tablename, con=conn, index=False, if_exists='append') print(f"成功插入 {len(valid_df)} 行") # 保存错误行到文件用于后续排查 if error_rows: error_df = pd.DataFrame([{**row, 'error_msg': msg, 'row_index': idx} for idx, msg, row in error_rows]) error_df.to_csv('error_rows.csv', index=False) print(f"发现 {len(error_df)} 行错误,已保存到error_rows.csv")
这种方法能最大程度保证插入效率,提前过滤错误行后可全量插入正常数据,同时保留错误行信息方便后续修复。
方案3:利用数据库原生批量导入工具(针对特定数据库)
如果使用MySQL/PostgreSQL等数据库,用原生批量导入工具配合错误处理会更高效:
MySQL示例(使用LOAD DATA LOCAL INFILE)
import pandas as pd from sqlalchemy import create_engine # 先将DataFrame保存为临时CSV temp_csv = 'temp_data.csv' df.to_csv(temp_csv, sep='|', index=False, header=False) engine = create_engine("mysql+pymysql://user:password@host/db") with engine.begin() as conn: # 执行LOAD DATA命令,忽略错误行(需数据库开启local_infile) conn.execute(f""" LOAD DATA LOCAL INFILE '{temp_csv}' INTO TABLE {tablename} FIELDS TERMINATED BY '|' OPTIONALLY ENCLOSED BY '"' ESCAPED BY '\\\\' IGNORE 1 LINES (@col1, @col2, @col3) SET id = NULLIF(@col1, ''), name = @col2, create_date = STR_TO_DATE(@col3, '%Y-%m-%d') ; """) # 查看警告信息获取错误行详情 warnings = conn.execute("SHOW WARNINGS").fetchall() if warnings: print(f"导入过程中出现警告:{warnings}")
PostgreSQL示例(使用COPY命令)
import pandas as pd from sqlalchemy import create_engine import io engine = create_engine("postgresql://user:password@host/db") with engine.begin() as conn: # 将DataFrame转为CSV格式的内存对象 output = io.StringIO() df.to_csv(output, sep='|', index=False, header=False) output.seek(0) # 使用COPY命令,忽略冲突或错误行 conn.execute(f""" COPY {tablename} FROM STDIN WITH (FORMAT csv, DELIMITER '|', QUOTE '"', ESCAPE '\\\\', HEADER false) ON CONFLICT DO NOTHING; """, data=output.getvalue())
原生工具的插入速度远快于pandas的to_sql,同时能通过数据库自身机制处理错误行(如忽略或记录)。
内容的提问来源于stack exchange,提问作者Southsun1
相关产品推荐
相关产品推荐

