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

如何用Apache Airflow动态Task Group处理百万级压缩文本文件

Apache Airflow 大ZIP文件并行处理解决方案

需求说明

  • 处理包含1~5000万行文本的ZIP文件,流程为:读取ZIP内文本→逐行转换→生成新ZIP
  • 后续需完成Postgres表更新、触发SFTP传输DAG
  • 单任务处理效率低,需通过并行拆分任务(例如1500万行拆分为6个任务,各处理250万行)
  • 需实现动态生成并行任务(替代固定偏移的硬编码方式),并完成结果合并

动态Task Group实现方案

1. 通用Chunk处理函数

将重复的任务逻辑抽象为通用函数,通过参数传递行号偏移:

import zipfile
import io
import os
from itertools import islice
from airflow.operators.python import PythonOperator

def apply_transformation(line):
    return f"{line}_NEW"

def process_chunk(**context):
    # 从DAG配置获取输入ZIP路径
    input_zip = context['dag_run'].conf.get("file_name")
    # 获取当前任务的处理行范围
    start_line = context["params"]["start_line"]
    end_line = context["params"]["end_line"]
    # 生成唯一临时输出文件,避免任务间冲突
    temp_output_path = f"/tmp/processed_chunk_{start_line}_{end_line}.txt"

    with zipfile.ZipFile(input_zip) as zf:
        for txt_name in zf.namelist():
            with io.TextIOWrapper(zf.open(txt_name), encoding="UTF-8") as input_fp:
                # 转换行号为文件迭代索引(文件对象从0开始计数)
                with open(temp_output_path, "w", encoding="UTF-8") as output_fp:
                    # islice(start, stop):跳过前start行,取到stop行(左闭右开)
                    for line in islice(input_fp, start_line - 1, end_line):
                        transformed_line = apply_transformation(line.rstrip("\n"))
                        output_fp.write(f"{transformed_line}\n")

2. 动态生成Task Group

在DAG中根据总行数和单任务处理量,循环生成并行任务:

from airflow import DAG
from airflow.utils.task_group import TaskGroup
from airflow.operators.python import PythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
import glob
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
}

def start_task(**context):
    print("Starting large file processing workflow...")

def merge_and_finalize(**context):
    # 1. 合并所有临时chunk文件
    temp_files = glob.glob("/tmp/processed_chunk_*.txt")
    # 按起始行号排序,保证输出顺序与原文件一致
    temp_files.sort(key=lambda x: int(x.split("_")[2]))
    
    # 2. 生成最终ZIP文件
    output_zip = "transformed_output.zip"
    with zipfile.ZipFile(output_zip, "w", zipfile.ZIP_DEFLATED) as zf:
        merged_content = ""
        for temp_file in temp_files:
            with open(temp_file, "r", encoding="UTF-8") as fp:
                merged_content += fp.read()
            os.remove(temp_file)  # 清理临时文件
        zf.writestr("transformed_data.txt", merged_content)
    
    # 3. 更新Postgres表(示例逻辑)
    # import psycopg2
    # conn = psycopg2.connect(host="postgres", dbname="mydb", user="user", password="pass")
    # with conn.cursor() as cur:
    #     cur.execute("UPDATE processing_status SET status = 'completed' WHERE file_name = %s", 
    #                 (context['dag_run'].conf.get("file_name"),))
    # conn.commit()
    # conn.close()
    
    # 4. 触发SFTP传输DAG
    TriggerDagRunOperator(
        task_id="trigger_sftp_transfer",
        trigger_dag_id="sftp_transfer_dag",
        conf={"file_path": output_zip}
    ).execute(context)

with DAG("large_file_processing",
         schedule_interval=None,
         default_args=default_args,
         catchup=False) as dag:

    start_op = PythonOperator(
        task_id='init_workflow',
        python_callable=start_task
    )

    # 配置并行参数:可从DAG conf传入或预计算
    total_lines = 15000000  # 示例总行数
    chunk_size = 2500000    # 每个任务处理行数
    num_chunks = total_lines // chunk_size
    # 处理余数,确保所有行都被覆盖
    if total_lines % chunk_size != 0:
        num_chunks += 1

    # 动态生成并行任务组
    with TaskGroup(group_id='parallel_chunk_processing') as parallel_tg:
        for chunk_idx in range(num_chunks):
            start_line = chunk_idx * chunk_size + 1
            end_line = min((chunk_idx + 1) * chunk_size, total_lines)
            PythonOperator(
                task_id=f"process_chunk_{chunk_idx + 1}",
                python_callable=process_chunk,
                params={"start_line": start_line, "end_line": end_line}
            )

    final_op = PythonOperator(
        task_id='finalize_processing',
        python_callable=merge_and_finalize
    )

    # 任务依赖
    start_op >> parallel_tg >> final_op

结果合并与后续操作

所有并行任务完成后,finalize_processing任务会:

  1. 收集并排序所有临时chunk文件,保证输出顺序与原文件一致
  2. 合并内容并压缩为最终ZIP
  3. 执行Postgres表更新逻辑
  4. 通过TriggerDagRunOperator触发SFTP传输DAG

其他并行处理方案

1. 动态任务映射(Airflow 2.2+)

无需手动循环生成Task,使用expand方法实现动态并行:

# 生成chunk参数列表
chunk_params = []
for chunk_idx in range(num_chunks):
    start_line = chunk_idx * chunk_size + 1
    end_line = min((chunk_idx + 1) * chunk_size, total_lines)
    chunk_params.append({"start_line": start_line, "end_line": end_line})

# 动态生成并行任务
process_tasks = PythonOperator.partial(
    task_id="process_chunk",
    python_callable=process_chunk
).expand(params=chunk_params)

# 依赖关系
start_op >> process_tasks >> final_op

2. Apache Spark分布式处理

对于5000万行级别的超大文件,使用Spark的分布式计算更高效:

  • 通过SparkSubmitOperator提交Spark作业,自动拆分文件并行处理
  • Spark支持直接读取ZIP文件,转换后写入新ZIP或存储系统

3. 预拆分文件

在初始化任务中先将原ZIP内的文本文件拆分为多个小文件,再分配给并行任务处理,避免行号偏移计算:

def split_large_file(**context):
    input_zip = context['dag_run'].conf.get("file_name")
    split_dir = "/tmp/split_files/"
    os.makedirs(split_dir, exist_ok=True)
    
    with zipfile.ZipFile(input_zip) as zf:
        for txt_name in zf.namelist():
            with io.TextIOWrapper(zf.open(txt_name), encoding="UTF-8") as fp:
                chunk_idx = 0
                while True:
                    chunk = list(islice(fp, chunk_size))
                    if not chunk:
                        break
                    with open(f"{split_dir}split_{chunk_idx}.txt", "w", encoding="UTF-8") as out_fp:
                        out_fp.writelines(chunk)
                    chunk_idx +=1

后续并行任务直接处理拆分后的小文件即可。


内容的提问来源于stack exchange,提问作者PROTOCOL

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 17:55:22