使用cx_Oracle编写可扩展INSERT语句实现DataFrame插入Oracle
解决方案:自动匹配列的DataFrame批量插入Oracle ODS表
方案一:基于cx_Oracle的原生批量插入(推荐,无连接问题)
你已通过cx_Oracle正常连接数据库,且能获取DataFrame与ODS表的共同列,接下来可通过动态生成INSERT语句+批量绑定变量实现可扩展插入,无需硬编码列名:
完整步骤代码
import pandas as pd import cx_Oracle import config import numpy as np # 1. 读取数据 df = pd.read_excel("Employee_data.xlsx") conn = None try: # 2. 建立数据库连接 conn = cx_Oracle.connect(config.username, config.password, config.dsn, encoding=config.encoding) cursor = conn.cursor() # 3. 获取ODS表的列名(用WHERE 1=0避免查询数据,提升效率) cursor.execute("SELECT * FROM ODSMGR.EMPLOYEE_TABLE WHERE 1=0") col_names = [desc[0] for desc in cursor.description] # 4. 获取DataFrame与表的共同列 common_cols = np.intersect1d(df.columns, col_names).tolist() if not common_cols: raise ValueError("DataFrame与ODS表无匹配列") # 5. 过滤DataFrame并调整列顺序与表一致,避免数据错位 df_filtered = df[common_cols] rows = [tuple(row) for row in df_filtered.itertuples(index=False, name=None)] # 6. 动态生成INSERT语句 cols_str = ", ".join(common_cols) placeholders = ", ".join([f":{i+1}" for i in range(len(common_cols))]) insert_sql = f"INSERT INTO ODSMGR.EMPLOYEE_TABLE ({cols_str}) VALUES ({placeholders})" # 7. 批量插入(executemany比单条插入效率提升显著) cursor.executemany(insert_sql, rows) conn.commit() print(f"成功插入{len(rows)}条数据") except cx_Oracle.Error as error: if conn: conn.rollback() print(f"数据库错误: {error}") finally: # 关闭游标和连接 if cursor: cursor.close() if conn: conn.close()
关键说明
- 用
WHERE 1=0获取表结构,避免查询实际数据,提升效率 - 强制对齐DataFrame与表的列顺序,防止因列顺序不一致导致的数据插入错误
- 使用
executemany批量插入,适配大数据量场景 - 动态生成SQL语句,完全无需硬编码列名,适配任意数量的匹配列
方案二:修复SQLAlchemy连接错误后使用to_sql插入
你遇到的ORA-12505错误是因为SQLAlchemy连接字符串格式错误,正确格式需指定服务名/SID,而非用?拼接。
修正后的连接与插入代码
from sqlalchemy import create_engine, types, inspect import pandas as pd import numpy as np df = pd.read_excel("Employee_data.xlsx") # 正确的连接字符串格式(二选一) # 方式1:使用服务名(推荐) engine = create_engine("oracle+cx_oracle://username:password@ODS-test.domain.com:1521/ODS-Test") # 方式2:使用SID # engine = create_engine("oracle+cx_oracle://username:password@ODS-test.domain.com:1521/?service_name=ODS-Test") # 获取ODS表列名 inspector = inspect(engine) col_info = inspector.get_columns("ODSMGR.EMPLOYEE_TABLE") col_names = [col["name"] for col in col_info] # 过滤共同列 common_cols = np.intersect1d(df.columns, col_names).tolist() df_filtered = df[common_cols] # 定义object类型列的长度(避免默认长度不足报错) dtyp = {} for c in df_filtered.columns[df_filtered.dtypes == 'object']: max_len = df_filtered[c].str.len().max() dtyp[c] = types.VARCHAR(max_len if not pd.isna(max_len) else 100) # 批量插入 df_filtered.to_sql( name="EMPLOYEE_TABLE", con=engine, schema="ODSMGR", # 指定schema,避免表名识别错误 dtype=dtyp, if_exists="append", index=False, # 不插入DataFrame索引列 chunksize=1000 # 分块插入,避免内存溢出 )
关键修正点
- 修正连接字符串格式,去掉
?,直接用host:port/service_name指定服务名 - 添加
schema="ODSMGR",确保SQLAlchemy正确定位目标表 - 增加
chunksize参数,处理大数据量时分块插入 - 过滤DataFrame到共同列后再插入,避免列不匹配错误
两种方案对比
| 方案 | 优点 | 缺点 |
|---|---|---|
| cx_Oracle原生 | 连接稳定,性能更高,可灵活控制插入逻辑 | 需要手动处理列顺序和绑定变量 |
| SQLAlchemy to_sql | 代码简洁,无需手动编写INSERT语句 | 依赖SQLAlchemy配置,连接问题排查稍复杂 |
内容的提问来源于stack exchange,提问作者Benjamin Diaz
相关产品推荐
相关产品推荐

