SQLAlchemy如何将查询返回的结果对象插入目标MySQL数据表
跨MySQL实例查询结果写入目标表实现方法
你当前持有的result是SQLAlchemy返回的本地游标结果对象,不属于任何一个MySQL服务端的资源,目标库连接无法直接读取本地Python变量里的结果,因此无法直接执行Insert into DB.TBL SELECT * FROM RESULT这类直接引用本地结果的SQL,可按数据量规模选以下两种实现方案:
方案1:百万行以内小数据量(最简实现)
直接用Pandas做中转,代码改动量最小,不需要手动拼接插入SQL:
import pandas as pd from sqlalchemy import create_engine try: # 初始化连接 engine_source = create_engine("源库实际连接字符串", pool_pre_ping=True) engine_dest = create_engine("目标库实际连接字符串", pool_pre_ping=True) # 读取源库查询结果到DataFrame source_query = '这里替换成你的实际SELECT查询语句' df = pd.read_sql(source_query, con=engine_source) print(f'EXTRACT COMPLETE,共读取{len(df)}行数据') # 批量写入目标表,效果等价于全量结果插入 df.to_sql( name='替换为目标表名', con=engine_dest, schema='替换为目标表所在的数据库名', if_exists='append', # 追加模式,不会覆盖目标表现有结构和数据 index=False, # 不写入Pandas自动生成的索引列 chunksize=10000 # 每1万行提交一次批量插入,避免单次提交数据量过大报错 ) print('INSERT COMPLETE') # 释放连接 engine_source.dispose() engine_dest.dispose() except Exception as e: print('task error: ' + str(e))
注意:使用该方法必须保证源查询返回的列顺序、字段类型和目标表完全一致,否则会出现字段错位、类型转换失败问题。
方案2:百万行以上大数据量(低内存占用高性能实现)
如果全量读入DataFrame会占满内存,就用服务端游标分批拉取、批量插入的方式,全程不需要把全量数据加载到内存:
from sqlalchemy import create_engine, text try: engine_source = create_engine("源库实际连接字符串", pool_pre_ping=True) engine_dest = create_engine("目标库实际连接字符串", pool_pre_ping=True) source_query = '这里替换成你的实际SELECT查询语句' target_table = '替换为目标库名.目标表名' batch_size = 10000 # 每次拉取/插入的批次大小 with engine_source.connect() as src_conn, engine_dest.connect() as dest_conn: # 开启服务端游标,不一次性加载全量结果 result_proxy = src_conn.execution_options(stream_results=True).execute(text(source_query)) # 获取列名,构造插入SQL columns = [col.name for col in result_proxy.cursor.description] insert_sql = f"INSERT INTO {target_table} ({','.join(columns)}) VALUES ({','.join([':%s' % i for i in range(len(columns))])})" while True: # 分批拉取数据 batch = result_proxy.fetchmany(batch_size) if not batch: break # 批量插入 dest_conn.execute(text(insert_sql), [dict(zip(columns, row)) for row in batch]) dest_conn.commit() print('INSERT COMPLETE') engine_source.dispose() engine_dest.dispose() except Exception as e: print('task error: ' + str(e))
注意事项
- 如果目标表有自增主键,要么源查询排除自增列由目标库自动生成,要么保证源数据主键值不重复,否则会触发主键冲突报错
- 超大数据量写入时,可以暂时关闭目标表的二级索引、外键约束,写入完成后再重建,能大幅提升写入速度
- 跨库写入前建议先拿100行测试数据做验证,确认字段映射关系正确后再跑全量任务,避免写入脏数据
内容的提问来源于stack exchange,提问作者kikee1222
相关产品推荐
相关产品推荐

