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

Airflow Python全局声明pandas DataFrame合并时df1为空报错求解

问题解决方案

根因说明

Airflow的每一个任务都会运行在独立的进程空间中,不同进程的全局变量完全隔离,无法共享状态。你在before_load任务中修改的全局df1仅存在于该任务的进程里,后续send_success_email任务启动时会加载全新的代码上下文,拿到的还是初始化时的空DataFrame,合并时找不到ATTRIBUTE字段就会抛出KeyError。

解决方案

使用Airflow官方提供的XCom组件实现跨任务数据传递,不要依赖全局变量传值,修改点如下:

  1. 删除原有全局df1 = pd.DataFrame()的声明
  2. 修改get_before_load_count函数,将计算结果推送到XCom
  3. 修改send_success_notification函数,从XCom拉取前置任务的计算结果

修改后核心代码

# 原有导入、默认参数、sf_hook等部分保留不变,仅修改两个函数部分即可
def get_before_load_count(**context):
    query1 = """SELECT 'att1' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_BEFORE_LOAD FROM DB.SCHEMA.TABLE1
               UNION ALL
               SELECT 'att2' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_BEFORE_LOAD FROM DB.SCHEMA.TABLE2
               UNION ALL
               SELECT 'att3' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_BEFORE_LOAD FROM DB.SCHEMA.TABLE3
               UNION ALL
               SELECT 'att4' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_BEFORE_LOAD FROM DB.SCHEMA.TABLE4
               UNION ALL
               SELECT 'att5' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_BEFORE_LOAD FROM DB.SCHEMA.TABLE5 """
    df1 = pd.read_sql(query1, sf_con)
    print("df1", df1)
    # 将结果转为可序列化格式推送到XCom
    context['ti'].xcom_push(key='before_load_count', value=df1.to_dict('records'))

def send_success_notification(**context):
    # 从XCom拉取前置任务的df1数据并转回DataFrame
    df1_records = context['ti'].xcom_pull(key='before_load_count', task_ids='before_load')
    df1 = pd.DataFrame(df1_records)
    
    query = """SELECT 'att1' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_AFTER_LOAD FROM DB.SCHEMA.TABLE1
               UNION ALL
               SELECT 'att2' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_AFTER_LOAD FROM DB.SCHEMA.TABLE2
               UNION ALL
               SELECT 'att3' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_AFTER_LOAD FROM DB.SCHEMA.TABLE3
               UNION ALL
               SELECT 'att4' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_AFTER_LOAD FROM DB.SCHEMA.TABLE4
               UNION ALL
               SELECT 'att5' AS ATTRIBUTE,
               COUNT(*) AS RECORD_COUNT_AFTER_LOAD FROM DB.SCHEMA.TABLE5"""
    df2 = pd.read_sql(query, sf_con)
    print("df2", df2)
    print("my df1", df1)
    df3 = pd.merge(df1, df2, on = "ATTRIBUTE")
    df = df3.reindex(["ATTRIBUTE","RECORD_COUNT_BEFORE_LOAD","RECORD_COUNT_AFTER_LOAD"], axis=1)
    print("df", df)
    html_table = df.to_html(index=False, justify='center')
    op = EmailOperator(
        task_id='success_email',
        to='xxx@example.com',
        subject='Email subject '+ date,
        html_content=" <p>Hi,<br><br>Process Completed<br><br> {}".format(html_table),
        dag=dag
    )
    op.execute(context)

注意事项

该场景下传递的数据量极小,使用XCom完全没有性能问题。如果后续需要传递更大体量的中间结果,可以考虑将结果写入临时表/对象存储,再在下游任务读取对应地址的数据即可。

内容的提问来源于stack exchange,提问作者Rajalakshmi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 15:57:03