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

如何通过Airflow并行运行Python差分进化脚本的50次迭代?

用Airflow并行执行差分进化脚本的50次迭代

需求背景

我有一个名为differential_evolution.py的Python差分进化脚本,单次迭代运行约40代。脚本已通过随机种子保证每次迭代结果独立,当前是串行执行50次迭代,耗时过长。希望通过Airflow重构DAG,让50次迭代并行执行,且每次迭代输出带序号的文件:abc_i.txt和xyz_i.csv(i为迭代序号)。

原脚本核心逻辑片段:

for iteration in range(50):
    seed = np.random.randint(0, 1000)  
    opt_obj = Optimizer()  
    solution = opt_obj.run_optimizer()  
    mape = opt_obj.calc_performance(solution.x)

当前DAG结构(串行):
start >> create_cluster >> differential_evolution.py >> delete_cluster >> end

期望DAG结构(并行):
start >> create_cluster >> [迭代1, 迭代2, ..., 迭代50] >> delete_cluster >> end


实现步骤

1. 修改Python脚本,支持参数化

将原脚本的循环逻辑拆分为单次迭代逻辑,接受外部传入的迭代序号,输出带序号的文件:

import numpy as np
import sys

# 从命令行获取迭代序号
iteration_idx = int(sys.argv[1])
# 生成随机种子(也可由Airflow传入,提升可控性)
seed = np.random.randint(0, 1000)

# 初始化优化器并绑定种子(假设Optimizer支持种子参数)
opt_obj = Optimizer(seed=seed)
# 若Optimizer通过方法设置种子,可替换为:opt_obj.set_seed(seed)

solution = opt_obj.run_optimizer()
mape = opt_obj.calc_performance(solution.x)

# 写入带序号的输出文件
with open(f'abc_{iteration_idx}.txt', 'w') as f:
    f.write(f"MAPE: {mape}\nSolution: {solution.x}")

with open(f'xyz_{iteration_idx}.csv', 'w') as f:
    # 示例:将解列表写入CSV
    f.write(','.join(map(str, solution.x)) + '\n')

2. 重构Airflow DAG,生成并行任务

在DAG脚本中通过循环创建50个并行任务,设置正确的依赖关系:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator
# 根据集群类型导入对应Operator,比如AWS EMR/GCP Dataproc的创建/删除Operator
from datetime import datetime

def run_single_iteration(iteration_idx):
    """调用修改后的差分进化脚本执行单次迭代"""
    import subprocess
    # 替换为脚本实际路径
    subprocess.run(
        ["python", "/opt/airflow/dags/scripts/differential_evolution.py", str(iteration_idx)],
        check=True
    )

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

with DAG(
    dag_id='parallel_differential_evolution',
    default_args=default_args,
    schedule_interval=None,  # 按需设置调度周期,比如@daily
    catchup=False
) as dag:
    # 起始/结束占位任务
    start = DummyOperator(task_id='start')
    end = DummyOperator(task_id='end')

    # 替换为实际的集群创建/删除Operator
    create_cluster = DummyOperator(task_id='create_cluster')
    delete_cluster = DummyOperator(task_id='delete_cluster')

    # 生成50个并行迭代任务
    iteration_tasks = []
    for i in range(1, 51):
        task = PythonOperator(
            task_id=f'de_iteration_{i}',
            python_callable=run_single_iteration,
            op_kwargs={'iteration_idx': i},
            # 按需添加资源限制参数,比如execution_timeout、pool等
        )
        iteration_tasks.append(task)

    # 设置依赖链
    start >> create_cluster >> iteration_tasks >> delete_cluster >> end

关键注意事项

  • 种子独立性:确保Optimizer类正确使用传入的随机种子,避免不同任务间的结果干扰。
  • 文件存储:若在分布式集群环境运行,需将输出文件写入共享存储(如S3、HDFS),避免本地文件无法跨节点访问。
  • 资源控制:根据Airflow集群的资源配置,通过DAG的max_active_runs或任务的pool参数限制并行数量,防止资源耗尽。
  • 集群适配:将示例中的DummyOperator替换为实际的集群操作Operator,比如AWS EMR的EmrCreateJobFlowOperator或GCP Dataproc的DataprocClusterCreateOperator。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 07:15:36