讲师中心 微信公众号
AI工具推荐 视频效率加速

如何在Airflow中安全高效地并行执行Pandas数据处理任务

千墨酱_1762

千墨酱_1762

发布时间:2026-05-13 16:03:03

|

511人浏览过

|

来源于php中文网

原创

如何在Airflow中安全高效地并行执行Pandas数据处理任务

Airflow不支持直接在XCom中传递大型DataFrame,因此需通过共享存储(如S3、本地挂载卷)或数据库中转实现任务间数据流转,再利用>>操作符显式声明依赖关系,确保清洗、转换等无依赖步骤可真正并行执行。

airflow不支持直接在xcom中传递大型dataframe,因此需通过共享存储(如s3、本地挂载卷)或数据库中转实现任务间数据流转,再利用>>操作符显式声明依赖关系,确保清洗、转换等无依赖步骤可真正并行执行。

在Airflow中对Pandas DataFrame函数进行“并行化”,关键不在于让pandas.apply()多线程运行(那是底层计算优化),而在于将逻辑解耦为独立、可调度、无状态的Airflow任务,并通过外部媒介交换数据——因为XCom设计初衷仅用于传递轻量元数据(如文件路径、记录数、状态码),而非GB级DataFrame。

你当前的data_transformation.py结构合理,但直接在DAG中串联调用并传递df对象会失败:Airflow Worker进程彼此隔离,df无法跨PythonOperator内存传递;且XCom默认有48KB大小限制(即使调大也严重拖慢性能并占用数据库)。

✅ 正确实践路径如下:

1. 重构为「输入–处理–输出」三段式模块

将每个函数改造为接受输入路径/表名、执行处理、写入明确输出路径/临时表:

# etl/transform.py
import pandas as pd
import os
from pathlib import Path

def get_df_from_db(output_path: str):
    """从DB读取原始数据,保存为Parquet(推荐)或CSV"""
    conn = ...  # 同原逻辑
    df = pd.read_sql("SELECT * FROM table", conn)
    # 推荐使用Parquet:列式、压缩、Schema保留
    df.to_parquet(output_path, index=False)
    return output_path  # 返回路径,供下游任务使用

def clean_colorTypes(input_path: str, output_path: str):
    df = pd.read_parquet(input_path)
    df['colorType'] = df['colorType'].apply(
        lambda x: x if x in ['Red', 'Blue', 'Orange'] else 'Others'
    )
    df = df[df['numOfButtons'] <= 2]
    df.to_parquet(output_path, index=False)

def add_yearsSinceManufactured(input_path: str, output_path: str):
    df = pd.read_parquet(input_path)
    df['yearsSinceManu'] = df['yearListed'] - df['yearManufactured']
    df.to_parquet(output_path, index=False)

def add_to_db(input_path: str):
    df = pd.read_parquet(input_path)
    # 使用SQLAlchemy engine写入目标表
    df.to_sql('cleaned', con=engine, if_exists='append', index=False)

2. 在DAG中声明并行任务与显式依赖

使用PythonOperator调用上述函数,并通过op_kwargs传入路径。注意:两个清洗任务共用同一输入路径,但各自写入不同中间文件,天然可并行:

# dags/my_etl_dag.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
from etl.transform import (
    get_df_from_db,
    clean_colorTypes,
    add_yearsSinceManufactured,
    add_to_db
)
import os

default_args = {
    'owner': 'data-engineer',
    'retries': 2,
    'retry_delay': timedelta(minutes=1),
    'execution_timeout': timedelta(minutes=15),  # 防止卡死
}

dag = DAG(
    'pandas_parallel_etl',
    default_args=default_args,
    schedule_interval='@daily',
    start_date=datetime(2026, 1, 1),
    catchup=False,
)

# Step 1: 提取原始数据 → 写入共享存储
extract_task = PythonOperator(
    task_id='extract_from_db',
    python_callable=get_df_from_db,
    op_kwargs={'output_path': '/shared/data/raw.parquet'},  # 绝对路径或S3 URI
    dag=dag,
)

# Step 2: 两个独立转换任务 → 并行执行
clean_task = PythonOperator(
    task_id='clean_colorTypes',
    python_callable=clean_colorTypes,
    op_kwargs={
        'input_path': '/shared/data/raw.parquet',
        'output_path': '/shared/data/cleaned_colors.parquet'
    },
    dag=dag,
)

years_task = PythonOperator(
    task_id='add_yearsSinceManufactured',
    python_callable=add_yearsSinceManufactured,
    op_kwargs={
        'input_path': '/shared/data/raw.parquet',
        'output_path': '/shared/data/enriched_years.parquet'
    },
    dag=dag,
)

# Step 3: 合并结果并入库(需先等待两个上游完成)
def merge_and_load(**context):
    # 读取两个中间结果,按主键/索引合并(示例:假设都含id列)
    df1 = pd.read_parquet('/shared/data/cleaned_colors.parquet')
    df2 = pd.read_parquet('/shared/data/enriched_years.parquet')
    merged = df1.merge(df2, on='id', how='inner')  # 或根据业务逻辑join
    add_to_db(merged)  # 或复用原add_to_db函数,传入merged

merge_task = PythonOperator(
    task_id='merge_and_load',
    python_callable=merge_and_load,
    dag=dag,
)

# 显式声明依赖:提取 → 并行清洗 → 合并入库
extract_task >> [clean_task, years_task] >> merge_task

⚠️ 关键注意事项

  • 路径必须绝对化:避免./data/,改用/opt/airflow/shared/(NFS挂载)或S3://bucket/path/(配合ObjectStoragePath);
  • 禁用GUI后端:若脚本中隐含绘图,开头加import matplotlib; matplotlib.use('Agg');
  • 异常必须捕获并抛出AirflowException:防止任务假死;
  • 大文件慎用pd.read_excel():优先转为Parquet/CSV,或用BashOperator + aws s3 cp预下载;
  • 不要在task内动态修改sys.path:模块应置于$AIRFLOW_HOME/dags/etl/并带__init__.py,用from etl.transform import xxx导入。

✅ 总结:Airflow中的“并行”本质是任务级并发,而非单任务内多线程。把DataFrame当作不可变的中间产物,用文件/数据库作为契约载体,才能构建健壮、可观测、可重试的生产级ETL流水线。

热门AI工具

更多
讯飞绘文

讯飞绘文是一款由科大讯飞推出的一站式 AIGC 内容运营平台。

UP简历
UP简历 Hot

一款AI办公效率工具,主要用于基于AI技术的免费在线简历制作工具,适合需要提升相关任务效率的用户。

WorkBuddy

一款AI办公效率工具,主要用于腾讯云推出的AI原生桌面智能体工作台,适合需要提升相关任务效率的用户。

二狗PPT
二狗PPT Hot

一款AI演示文稿工具,主要用于专为中式职场打造的AI PPT生成工具,适合需要提升相关任务效率的用户。

豆包大模型

豆包大模型是一款由字节跳动推出的企业级大语言模型服务平台。

蛙蛙写作

一款AI论文写作工具,主要用于超级AI智能写作助手,适合需要提升相关任务效率的用户。

音述AI
音述AI Hot

一款AI音频处理工具,主要用于音述AI是一个以“用声音述说故事”为核心的 AI 音乐创作与声音分享社区,适合需要提升相关任务效率的用户。

LibLibAI
LibLibAI Hot

一款AI视频创作工具,主要用于国内领先的AI创意平台,以海量模型、低门槛操作与“创作-分享-商业化”生态,让小白与专业创作者都能高效实现图文乃至视频创意表达,适合需要提升相关任务效率的用户。

DeepSeek

DeepSeek是一款面向对话、写作、编程和推理场景的AI大模型工具。

相关专题

更多
python打包成可执行文件
python打包成可执行文件

本专题为大家带来python打包成可执行文件相关的文章,大家可以免费的下载体验。

1691

2023.07.20

python能做什么
python能做什么

python能做的有:可用于开发基于控制台的应用程序、多媒体部分开发、用于开发基于Web的应用程序、使用python处理数据、系统编程等等。本专题为大家提供python相关的各种文章、以及下载和课程。

4284

2023.07.25

format在python中的用法
format在python中的用法

Python中的format是一种字符串格式化方法,用于将变量或值插入到字符串中的占位符位置。通过format方法,我们可以动态地构建字符串,使其包含不同值。php中文网给大家带来了相关的教程以及文章,欢迎大家前来阅读学习。

1689

2023.07.31

python教程
python教程

Python已成为一门网红语言,即使是在非编程开发者当中,也掀起了一股学习的热潮。本专题为大家带来python教程的相关文章,大家可以免费体验学习。

24997

2023.08.03

python环境变量的配置
python环境变量的配置

Python是一种流行的编程语言,被广泛用于软件开发、数据分析和科学计算等领域。在安装Python之后,我们需要配置环境变量,以便在任何位置都能够访问Python的可执行文件。php中文网给大家带来了相关的教程以及文章,欢迎大家前来学习阅读。

3047

2023.08.04

python eval
python eval

eval函数是Python中一个非常强大的函数,它可以将字符串作为Python代码进行执行,实现动态编程的效果。然而,由于其潜在的安全风险和性能问题,需要谨慎使用。php中文网给大家带来了相关的教程以及文章,欢迎大家前来学习阅读。

3067

2023.08.04

scratch和python区别
scratch和python区别

scratch和python的区别:1、scratch是一种专为初学者设计的图形化编程语言,python是一种文本编程语言;2、scratch使用的是基于积木的编程语法,python采用更加传统的文本编程语法等等。本专题为大家提供scratch和python相关的文章、下载、课程内容,供大家免费下载体验。

1163

2023.08.11

python合并两个列表
python合并两个列表

Python是一种强大的编程语言,具有许多方便的功能和工具。在Python中,有多种方法可以合并两个列表。php中文网给大家带来了相关的教程以及文章,欢迎大家前来学习阅读。

596

2023.08.10

PixTV官网入口地址合集
PixTV官网入口地址合集

本专题汇总了 PixTV AI 一站式视频创作平台的官方入口与使用教程。无需下载软件,浏览器直接访问即可使用。平台将剧本、图像、视频、声音与剪辑整合在“无限画布”中,接入 GPT Image 2.5、Seedance 2.5 等头部模型。本专题整理了从新建画布、角色锚定、分镜拆分到视频生成与导出的完整操作指南,助你快速上手 AI 短剧与漫剧创作。

20

2026.10.10

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn