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

能否在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中按以下流程执行:

  1. 执行get_df_from_db获取数据;
  2. 并行执行clean_colorTypes和add_yearsSinceManufactured(二者无依赖关系);
  3. 完成前两步后执行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 22:27:35