能否在Airflow中并行执行Pandas DataFrame数据处理函数?
Airflow学术项目流程可行性分析与修正方案
用户问题背景
我正在开展一个需使用Airflow的学术项目,目前仅处理一张数据库表,不确定当前思路是否正确。我有一个名为data_transformation.py的文件,内容如下:
def get_df_from_db(): conn = pymysql.connect(host='HOST', user='USER', password='PW', database='DB') query = "SELECT * FROM table" df = pd.read_sql(query, conn) df.reset_index(inplace=True) df.drop('index', axis=1, inplace=True) connection.close() return df def clean_colorTypes(df): df['colorType'] = df['colorType'].apply(lambda x: x if x in ['Red', 'Blue', 'Orange'] else 'Others') df = df[df['numOfButtons'] <= 2] return df def add_yearsSinceManufactured(df): df['yearsSinceManu'] = df['yearListed'] - df['yearManufactured'] return df def add_to_db(df): #create connection to mysql using create_engine df.to_sql(name = 'cleaned', con = conn, if_exists = 'append', index = False) return df
计划在Airflow DAG中按以下流程执行:
- 执行
get_df_from_db获取数据; - 并行执行
clean_colorTypes和add_yearsSinceManufactured(二者无依赖关系); - 完成前两步后执行
add_to_db写入数据库,但由于每个函数都传递DataFrame,不确定该方案是否可行。
核心问题与修正方案
1. Airflow任务无法直接传递内存中的DataFrame
Airflow的每个任务都是独立进程(甚至跨机器)运行,任务间不能直接传递内存里的DataFrame。必须通过中间存储传递数据,比如:
- 将第一步获取的DataFrame写入临时存储(本地Parquet文件、数据库临时表等);
- 后续任务从临时存储读取数据,处理后再写回;
- 最终任务合并处理结果后写入目标表。
2. 并行任务的逻辑冲突
你计划并行的两个函数存在数据不一致风险:
clean_colorTypes会过滤掉numOfButtons > 2的行,而add_yearsSinceManufactured是在全量数据上新增列,两者输出的DataFrame行数/结构不匹配,后续合并会出错。- 修正方向:要么先执行过滤再计算新增列(串行),要么让两个并行任务处理同一份原始数据,最终按主键合并结果(如果业务允许)。
3. 代码中的基础错误
get_df_from_db中,连接变量是conn,关闭时误用connection.close(),变量名不匹配会报错;add_to_db中未定义conn,需通过create_engine创建合法连接;reset_index后立即drop索引列是多余操作,pd.read_sql默认不会把查询结果的索引存入DataFrame,可删除该部分代码。
修正后的实现示例
调整数据传递逻辑
新增临时数据读写函数:
import pandas as pd def save_temp_df(df, file_path): # 用Parquet格式存储,比CSV更高效且保留数据类型 df.to_parquet(file_path, index=False) def load_temp_df(file_path): return pd.read_parquet(file_path)
修正核心函数错误
import pymysql from sqlalchemy import create_engine def get_df_from_db(): conn = pymysql.connect(host='HOST', user='USER', password='PW', database='DB') query = "SELECT * FROM table" df = pd.read_sql(query, conn) conn.close() # 修正变量名 return df def add_to_db(df): # 创建SQLAlchemy引擎 engine = create_engine('mysql+pymysql://USER:PW@HOST/DB') df.to_sql(name='cleaned', con=engine, if_exists='append', index=False) return df
Airflow DAG结构示例
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import data_transformation as dt # 临时文件路径(根据Airflow部署环境调整,确保所有任务节点可访问) TEMP_DATA_PATH = "/tmp/raw_data.parquet" CLEANED_DATA_PATH = "/tmp/cleaned_data.parquet" CALCULATED_DATA_PATH = "/tmp/calculated_data.parquet" with DAG( dag_id="data_processing_dag", start_date=datetime(2024,1,1), schedule_interval="@daily", catchup=False ) as dag: # 任务1:获取原始数据并保存到临时存储 get_raw_data = PythonOperator( task_id="get_raw_data", python_callable=lambda: dt.save_temp_df(dt.get_df_from_db(), TEMP_DATA_PATH) ) # 任务2a:清洗数据 clean_data = PythonOperator( task_id="clean_data", python_callable=lambda: dt.save_temp_df(dt.clean_colorTypes(dt.load_temp_df(TEMP_DATA_PATH)), CLEANED_DATA_PATH) ) # 任务2b:计算新增列 calculate_years = PythonOperator( task_id="calculate_years", python_callable=lambda: dt.save_temp_df(dt.add_yearsSinceManufactured(dt.load_temp_df(TEMP_DATA_PATH)), CALCULATED_DATA_PATH) ) # 任务3:合并处理结果并写入数据库 def merge_and_write(): cleaned_df = dt.load_temp_df(CLEANED_DATA_PATH) calculated_df = dt.load_temp_df(CALCULATED_DATA_PATH) # 按主键合并(假设主键是id) final_df = cleaned_df.merge(calculated_df[['id', 'yearsSinceManu']], on='id', how='left') dt.add_to_db(final_df) merge_and_write_db = PythonOperator( task_id="merge_and_write_db", python_callable=merge_and_write ) # 设置任务依赖 get_raw_data >> [clean_data, calculate_years] >> merge_and_write_db
总结
你的提取-处理-加载的核心思路是可行的,但需要解决Airflow任务间的数据传递问题,修正代码中的基础错误,并调整并行任务的逻辑以避免数据不一致。对于学术项目,本地Parquet文件是简单高效的中间存储方案;如果后续扩展到生产环境,可替换为分布式存储或数据库临时表。
内容的提问来源于stack exchange,提问作者vvv-1234
相关产品推荐
相关产品推荐

