Airflow Python全局声明pandas DataFrame合并时df1为空报错求解
问题解决方案
根因说明
Airflow的每一个任务都会运行在独立的进程空间中,不同进程的全局变量完全隔离,无法共享状态。你在before_load任务中修改的全局df1仅存在于该任务的进程里,后续send_success_email任务启动时会加载全新的代码上下文,拿到的还是初始化时的空DataFrame,合并时找不到ATTRIBUTE字段就会抛出KeyError。
解决方案
使用Airflow官方提供的XCom组件实现跨任务数据传递,不要依赖全局变量传值,修改点如下:
- 删除原有全局
df1 = pd.DataFrame()的声明 - 修改
get_before_load_count函数,将计算结果推送到XCom - 修改
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
相关产品推荐
相关产品推荐

