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

MWAA触发DAG时遭遇DagRunAlreadyExists错误求助

MWAA触发DAG时偶现DagRunAlreadyExists错误处理

问题概述

触发MWAA中的DAG时偶尔出现核心错误:

airflow.exceptions.DagRunAlreadyExists: 指定dag_id、execution_date和run_id的DAG实例已存在

相关代码片段

def should_retry(response) -> bool:
    """
    If `true`, then retry the Airflow job. Else, fail.
    """
    wait_between_s3_requests()
    response_data = decode_response_obj(response)
    print(response_data)
    if "externally triggered: True" in response_data["stdout"].decode("utf8") and response.status_code == 200:
        print("here")
        send_slack_message(f"here")
        return False
    if response is not None and response.status_code != 200:
        return True
    if response is not None and response.status_code == 200:
        response_data = decode_response_obj(response)
        if response_data["stderr"]:
            return True
        else:
            return False
    return False


def _trigger_dag(dag_name: str, *, conf: str, airflow_env: str) -> http.client.HTTPSConnection:
    """Trigger a DAG run in an Airflow environment."""
    mwaa_cli_token = boto3.client("mwaa").create_cli_token(Name=airflow_env)
    connection = http.client.HTTPSConnection(mwaa_cli_token["WebServe

完整错误信息

Trigger dag  with error: b'/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/configuration.py:357 DeprecationWarning: The processor_poll_interval option in [scheduler] has been renamed to scheduler_idle_sleep_time - the old setting has been used, but please update your config.\nTraceback (most recent call last):\n  File "/usr/local/airflow/.local/bin/airflow", line 8, in <module>\n    sys.exit(main())\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/__main__.py", line 48, in main\n    args.func(args)\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/cli/cli_parser.py", line 48, in command\n    return func(*args, **kwargs)\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/utils/cli.py", line 92, in wrapper\n    return f(*args, **kwargs)\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/cli/commands/dag_command.py", line 138, in dag_trigger\n    dag_id=args.dag_id, run_id=args.run_id, conf=args.conf, execution_date=args.exec_date\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/api/client/local_client.py", line 30, in trigger_dag\n    dag_id=dag_id, run_id=run_id, conf=conf, execution_date=execution_date\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/api/common/experimental/trigger_dag.py", line 125, in trigger_dag\n    replace_microseconds=replace_microseconds,\n  File "/usr/local/airflow/.local/lib/python3.7/site-packages/airflow/api/common/experimental/trigger_dag.py", line 75, in _trigger_dag\n    f"A Dag Run already exists for dag id {dag_id} at {execution_date} with run id {run_id}"\nairflow.exceptions.DagRunAlreadyExists: A Dag Run already exists for dag id ml_product at 2023-01-14 00:08:19+00:00 with run id manual__2023-01-14T00:08:19+00:00\n'

错误原因

该错误是由于触发请求中使用的dag_id、execution_date和run_id组合已存在对应的DAG运行实例。结合代码重试逻辑来看,大概率是第一次触发成功后,重试机制未正确识别成功状态,重复发起了相同参数的请求,导致冲突。此外,_trigger_dag函数未完整实现,可能未为每次触发生成唯一的run_id。

解决方法

1. 生成唯一Run ID

每次触发DAG时生成唯一标识,避免重复:

import uuid
run_id = f"manual__{uuid.uuid4()}"

将该run_id传入触发命令,确保每次请求的标识唯一。

2. 优化重试判断逻辑

在should_retry函数中增加对DagRunAlreadyExists错误的检测,遇到该错误直接停止重试:

def should_retry(response) -> bool:
    wait_between_s3_requests()
    response_data = decode_response_obj(response)
    stdout_str = response_data["stdout"].decode("utf8")
    stderr_str = response_data["stderr"].decode("utf8") if response_data.get("stderr") else ""
    
    if "externally triggered: True" in stdout_str and response.status_code == 200:
        print("here")
        send_slack_message(f"here")
        return False
    # 检测到DAG已存在错误时不重试
    if "DagRunAlreadyExists" in stderr_str:
        return False
    if response is not None and response.status_code != 200:
        return True
    if response is not None and response.status_code == 200:
        return bool(stderr_str)
    return False

3. 完善_trigger_dag函数实现

确保构造请求时正确传递唯一run_id,完整示例如下:

def _trigger_dag(dag_name: str, *, conf: str, airflow_env: str):
    import uuid
    mwaa_cli_token = boto3.client("mwaa").create_cli_token(Name=airflow_env)
    connection = http.client.HTTPSConnection(mwaa_cli_token["WebServerHostname"])
    
    # 生成唯一run_id
    run_id = f"manual__{uuid.uuid4()}"
    command = f"dags trigger {dag_name} --conf '{conf}' --run-id '{run_id}'"
    headers = {
        'Authorization': f'Bearer {mwaa_cli_token["CliToken"]}',
        'Content-Type': 'text/plain'
    }
    
    connection.request("POST", "/aws_mwaa/cli", command, headers)
    response = connection.getresponse()
    return response

4. 检查DAG配置

确保目标DAG的catchup设置为False,避免自动生成重复的execution_date:

default_args = {
    'catchup': False
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 14:01:30