Airflow中pandas df.to_sql写入MySQL缺必填参数con报错如何解决
报错成因
- 作用域异常:你定义的
bq_extraction函数没有返回生成的DataFrame对象,执行df.to_sql时操作的df并非Pandas DataFrame实例,该同名对象的to_sql方法参数列表与Pandas不一致,触发参数缺失报错。 - 版本兼容问题:使用的Pandas版本低于1.3.0且SqlAlchemy版本高于2.0时,旧版Pandas的
to_sql方法无法正常识别高版本SqlAlchemy的连接对象,会误判con参数未传入。
解决方案
- 修正DataFrame作用域,给
bq_extraction函数添加返回值,确保写入操作的对象是合法的Pandas DataFrame:
def bq_extraction(): source_table = 'analytics.buyer_data' # Bigquery table destination_table = 'marketing.buyer_data' # Mysql Table sql_query = f"""SELECT email, count_products FROM {source_table} """ df = bq_client.query(query=sql_query, job_config=bq_query_config).to_dataframe() df.columns = ['email', 'count_wishlist_products'] print(df.shape) return df # 新增返回语句
- 显式指定
to_sql的所有必填参数,避免版本兼容导致的参数识别异常,同时拆分库名与表名通过schema参数指定库名,适配MySQL的写入规则:
# 先调用函数获取合法DataFrame df = bq_extraction() hook = MySqlHook('connection_id') engine = hook.get_sqlalchemy_engine() df.to_sql( name='buyers_data', schema='marketing', con=engine, if_exists='replace', # 可根据需求替换为append/fail index=False # 禁止写入DataFrame索引列,无需可删除 )
- 若确认是版本兼容问题,可通过如下命令调整依赖版本:
# 方案1:降级SqlAlchemy到1.4版本 pip install sqlalchemy==1.4.48 # 方案2:升级Pandas到1.3.0及以上版本 pip install pandas>=1.3.0
内容的提问来源于stack exchange,提问作者pythonlearner
相关产品推荐
相关产品推荐

