如何使用Python将大容量pandas DataFrame插入Hadoop表
报错原因及修复
报错的直接原因有两个:
executemany无法直接识别pandas.DataFrame类型的入参,要求传入的参数是可迭代的逐行数据序列(比如嵌套列表、元组列表)- 部分Hadoop ODBC驱动不支持
?作为参数占位符,需要替换为%s
修复后的游标插入代码:
with pyodbc.connect("DSN=hadoop",autocommit=True) as conx: cursor = conx.cursor() # 先提取目标字段转为逐行列表序列,占位符替换为%s适配驱动 insert_data = df[["keyword","category_l1","category_l2","brand","ordercode","sku","snp_subclass"]].values.tolist() cursor.executemany("INSERT INTO ast_labs_t.dcs_search_keywords_nlp_results_test (keyword,category_l1,category_l2,brand,ordercode,sku,snp_subclass) VALUES(%s,%s,%s,%s,%s,%s,%s)", insert_data)
注意:如果目标表的date字段不允许为空且没有默认值,需要把date字段也加入插入字段列表和取值逻辑中。
更高效率的插入方案
47000行属于中等数据量,游标插入的IO开销很高,推荐两种更优方案:
方案1:使用pandas to_sql批量插入(开发成本最低)
pandas的to_sql方法内置了批量插入优化,比手动写executemany效率高3~5倍,需要搭配sqlalchemy构造连接引擎:
from sqlalchemy import create_engine, String # 构造Hadoop连接引擎,可根据你的ODBC配置调整连接串 engine = create_engine("hive+pyodbc:///?DSN=hadoop") # 指定所有字段为String类型,匹配目标表结构 dtype = {col: String() for col in df.columns} # 批量插入,chunksize控制单次提交行数,避免请求超时 df.to_sql( name="dcs_search_keywords_nlp_results_test", schema="ast_labs_t", con=engine, if_exists="append", index=False, dtype=dtype, chunksize=1000 )
方案2:LOAD DATA文件加载(性能最高)
针对Hadoop生态表,最优的写入方式是直接加载数据文件,性能比ODBC插入高10倍以上:
- 把DataFrame导出为无表头的本地CSV文件:
df.to_csv("insert_data.csv", header=False, index=False, encoding="utf-8")
- 将CSV上传至HDFS路径
- 执行Hive/Spark SQL的LOAD语句直接加载数据:
LOAD DATA INPATH '/你的hdfs路径/insert_data.csv' INTO TABLE ast_labs_t.dcs_search_keywords_nlp_results_test
内容的提问来源于stack exchange,提问作者Akhilesh Pandey
相关产品推荐
相关产品推荐

