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
相关产品推荐
相关产品推荐

