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

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

千丽大大_8619

千丽大大_8619

发布时间:2026-05-13 14:46:04

|

670人浏览过

|

来源于php中文网

原创

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

本文详解如何在 airflow 中实现 pandas dataframe 函数的逻辑并行化,重点解决跨任务数据传递、内存与性能瓶颈、依赖编排等核心问题,并提供符合工程规范的替代方案。

本文详解如何在 airflow 中实现 pandas dataframe 函数的逻辑并行化,重点解决跨任务数据传递、内存与性能瓶颈、依赖编排等核心问题,并提供符合工程规范的替代方案。

在 Airflow 中直接“并行执行多个 Pandas 函数并共享 DataFrame”看似合理(如 clean_colorTypes 和 add_yearsSinceManufactured 同时运行),但该设计违背 Airflow 的核心范式,且存在严重隐患。根本原因在于:Airflow 的任务(Task)是进程隔离、状态无共享的执行单元,无法像本地脚本那样通过变量 df 在函数间直接传递——尤其当任务被调度到不同 Worker 节点时,内存中的 DataFrame 完全不可见。

❌ 为什么不能直接传 DataFrame?XCom 不是解决方案

你可能想到用 Airflow 的 XCom(Cross-Communication)机制传递 df,例如:

# 错误示范:禁止在 XCom 中传递大型 DataFrame
def task1(**context):
    df = get_df_from_db()
    context['task_instance'].xcom_push(key='raw_df', value=df)  # ⚠️ 内存爆炸!序列化开销巨大!

def task2(**context):
    df = context['task_instance'].xcom_pull(key='raw_df')  # ⚠️ 反序列化慢 + 占用 DB 存储

这会导致:

  • 内存溢出:Pandas DataFrame 序列化为 JSON/Pickle 后体积激增,XCom 默认后端(SQLite/PostgreSQL)会迅速撑爆;
  • 性能崩溃:10MB DataFrame 序列化+传输+反序列化可能耗时数秒至分钟,完全抵消并行收益;
  • 可扩展性归零:无法支持 >100MB 数据,更遑论生产级 ETL。

✅ 正确路径:解耦计算与数据,用存储代替传递

Airflow 的最佳实践是 “Orchestration, not Computation” —— 它负责调度与协调,而非执行重计算。因此,应将数据持久化为中间产物,让并行任务各自读取、独立处理、写回新结果:

✅ 推荐方案一:基于文件存储的轻量并行(推荐初学者 & 中小数据)

使用本地 NFS、S3 或 MinIO 存储 Parquet(高效列式格式),配合 ObjectStoragePath(Airflow 2.7+ 新特性)实现云中立路径抽象:

# dags/etl_pipeline.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.models import Variable
from datetime import datetime, timedelta
import pandas as pd
from pathlib import Path

# 使用 ObjectStoragePath 抽象(兼容 S3/GCS/Azure)
BASE_PATH = "s3://my-airflow-bucket/etl-staging/"

def fetch_and_save_raw_data(**context):
    df = get_df_from_db()  # 你的原函数(需修复 connection.close() → conn.close())
    # 保存为 Parquet(压缩 + 列式读取快)
    output_path = f"{BASE_PATH}raw_{context['ts_nodash']}.parquet"
    df.to_parquet(output_path, index=False)
    context['task_instance'].xcom_push(key='raw_path', value=output_path)

def clean_color_types(**context):
    input_path = context['task_instance'].xcom_pull(key='raw_path')
    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]
    output_path = input_path.replace("raw_", "cleaned_color_")
    df.to_parquet(output_path, index=False)
    context['task_instance'].xcom_push(key='color_path', value=output_path)

def add_years_since_manu(**context):
    input_path = context['task_instance'].xcom_pull(key='raw_path')
    df = pd.read_parquet(input_path)
    df['yearsSinceManu'] = df['yearListed'] - df['yearManufactured']
    output_path = input_path.replace("raw_", "enriched_year_")
    df.to_parquet(output_path, index=False)
    context['task_instance'].xcom_push(key='year_path', value=output_path)

def merge_and_load_to_db(**context):
    color_path = context['task_instance'].xcom_pull(key='color_path')
    year_path = context['task_instance'].xcom_pull(key='year_path')

    df_color = pd.read_parquet(color_path)
    df_year = pd.read_parquet(year_path)

    # 基于主键合并(假设 'id' 是唯一标识)
    merged = pd.merge(df_color, df_year[['id', 'yearsSinceManu']], on='id', how='left')
    # 写入目标库
    merged.to_sql('cleaned', con=conn, if_exists='append', index=False)

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

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

    extract = PythonOperator(
        task_id='fetch_raw_data',
        python_callable=fetch_and_save_raw_data,
        provide_context=True
    )

    clean = PythonOperator(
        task_id='clean_color_types',
        python_callable=clean_color_types,
        provide_context=True
    )

    enrich = PythonOperator(
        task_id='add_years_since_manu',
        python_callable=add_years_since_manu,
        provide_context=True
    )

    load = PythonOperator(
        task_id='merge_and_load_to_db',
        python_callable=merge_and_load_to_db,
        provide_context=True
    )

    # 显式声明依赖:clean 和 enrich 并行,均依赖 extract;load 依赖二者
    extract >> [clean, enrich] >> load

? 关键优势:

  • ✅ 真正并行:clean 与 enrich 任务由 Airflow Worker 并发拉起,互不阻塞;
  • ✅ 零 XCom 数据传输:仅传递轻量字符串路径(<1KB),无序列化开销;
  • ✅ 故障隔离:任一任务失败不影响另一分支,便于重试;
  • ✅ 可审计:所有中间 Parquet 文件留存,支持人工校验与调试。

✅ 推荐方案二:数据库内计算(推荐生产环境 & 大数据)

若数据已在 MySQL 中,优先将清洗逻辑下推至 SQL 层,利用数据库引擎并行能力,彻底规避 Pandas 瓶颈:

-- 创建清洗后视图或物化表(MySQL 8.0+ 支持 CTE)
CREATE TABLE cleaned AS
SELECT 
  *,
  CASE WHEN colorType IN ('Red','Blue','Orange') THEN colorType ELSE 'Others' END AS colorType_clean,
  yearListed - yearManufactured AS yearsSinceManu
FROM table 
WHERE numOfButtons <= 2;

Airflow DAG 仅需两个轻量任务:

  • BashOperator 执行 mysql -e "source /sql/clean.sql"
  • SqlSensor 校验 cleaned 表行数是否 > 0

此方案吞吐量提升 10–100 倍,且资源消耗趋近于零。

⚠️ 必须规避的陷阱(附修复代码)

问题 风险 修复方式
connection.close() 拼写错误(原文为 connection.close(),但变量名是 conn) 连接泄漏 → 数据库连接池耗尽 改为 conn.close()
pd.read_sql() 无 chunksize 大表加载内存溢出 对 >100 万行表,改用 pd.read_sql(query, conn, chunksize=10000) + pd.concat()
Matplotlib 未禁用 GUI 后端 任务卡在 “running” 状态 在每个 task 函数开头加 import matplotlib; matplotlib.use('Agg')
缺少超时与异常捕获 任务假死,监控失效 default_args 中设 execution_timeout,所有 I/O 包裹 try/except 并 raise AirflowException

总结:Airflow + Pandas 的黄金法则

  • Never pass DataFrame via XCom —— 用存储路径(S3/DB/File)替代;
  • Parallelize at the task level, not function level —— >> [taskA, taskB] 是正确语法;
  • Push down computation when possible —— SQL > Pandas > Spark for structured tabular ETL;
  • Treat Airflow as conductor, not orchestra —— 它调度 Pandas,但从不运行 Pandas 的核心计算。

遵循以上原则,你不仅能安全实现并行清洗,更能构建出可监控、可重试、可扩展的企业级数据管道。

热门AI工具

更多
豆包大模型

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

超级简历WonderCV

一款AI办公效率工具,主要用于免费求职简历模版下载制作,应届生职场人必备简历制作神器,适合需要提升相关任务效率的用户。

DeepSeek

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

立刻MV
立刻MV Hot

立刻MV是一款AI文本写作工具,AI 音乐视频(MV)创作工具。

UpDream
UpDream Hot

一款AI视频创作工具,主要用于哔哩哔哩推出的自研AI视频创作工具,适合需要提升相关任务效率的用户。

WorkBuddy

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

Seko
Seko Hot

一款AI视频创作工具,主要用于商汤科技推出的创编一体的AI短视频创作Agent,适合需要提升相关任务效率的用户。

蛙蛙写作

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

Laper
Laper Hot

Laper是专为编剧、导演和制片人推出的 AI 原生剧本创作工具。

相关专题

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

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

1671

2023.07.20

python能做什么
python能做什么

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

4184

2023.07.25

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

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

1669

2023.07.31

python教程
python教程

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

24317

2023.08.03

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

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

2967

2023.08.04

python eval
python eval

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

3007

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

FrankenPHP集成Laravel详细教程
FrankenPHP集成Laravel详细教程

本专题提供FrankenPHP集成Laravel的详细配置指南,全面解析运行原理、开发环境搭建、Caddyfile配置、Octane工作模式、数据库连接、队列任务、定时任务和生产环境优化,解决部署过程中常见的报错与兼容性问题。

0

2026.10.08

热门下载

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

精品课程

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

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