如何用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任务会:
- 收集并排序所有临时chunk文件,保证输出顺序与原文件一致
- 合并内容并压缩为最终ZIP
- 执行Postgres表更新逻辑
- 通过
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
相关产品推荐
相关产品推荐

