Airflow中合并Pandas DataFrame报KeyError,本地运行正常求助
问题描述
我有两个DAG分别向Airflow的Xcom发送JSON数据,最后一个DAG将这两个JSON转换为Pandas DataFrame并执行合并操作。本地运行代码正常,但在Airflow中执行时出现KeyError: 'table_name'错误。
本地测试代码(模拟Xcom数据)
# Snowflake (Xcom) snowflake_dict = [{"table_name":"leads","last_update_date_snow":1699269594000},{"table_name":"loan","last_update_date_snow":1698966458752}] # mySQL (Xcom) mysql_dict = [{"table_name":"leads","last_update_date_my":1699341986000},{"table_name":"loan","last_update_date_my":1699275608110}] df_snowflake = pd.DataFrame(snowflake_dict) df_mysql = pd.DataFrame(mysql_dict) merge_df = pd.merge(df_snowflake, df_mysql, on='table_name', how='inner')
DAG任务定义
merge_last_update = PythonOperator( task_id='merge', python_callable=merge_last_update, provide_context=True, dag=dag )
报错信息
File "/usr/local/airflow/dags/resources/Snowflake_sync/merge_last_update.py", line 31, in merge_last_update
merge_df = pd.merge(df_snowflake, df_mysql, on='table_name', how='inner')
...
KeyError: 'table_name'
可能的原因及解决方案
- Xcom数据获取错误:检查
merge_last_update函数中xcom_pull的参数,确认task_id、dag_id是否正确,拉取到的数据结构是否和本地测试的列表+字典格式一致。若拉取到单个字典或空值,会导致DataFrame缺失table_name列。 - 序列化/反序列化异常:Airflow默认用JSON序列化Xcom数据,若推送的数据包含特殊类型(如datetime),序列化后可能丢失结构。推送Xcom时可显式指定
serialize_json=True,拉取后先打印数据内容确认结构。 - 上下文传递问题:即使设置了
provide_context=True,也要确保函数中通过ti.xcom_pull()正确获取数据,而非直接使用本地测试变量。 - DataFrame构造前校验:在构造DataFrame前添加日志,打印从Xcom拉取的
snowflake_dict和mysql_dict,确认是否包含table_name键,数据是否完整。
内容的提问来源于stack exchange,提问作者ornachshon
相关产品推荐
相关产品推荐

