You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 14:19:55