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

从MWAA环境获取DAG详情遇阻,Boto3代码无返回结果求助

解决MWAA获取DAG调度频率的问题

核心问题

你当前使用的get_environment API不返回任何DAG相关元数据,DagProcessingConfigs字段本身就不存在于该API的响应结构中——这个API仅用于获取MWAA环境的基础配置(如实例类型、Airflow版本、存储位置等),无法拿到DAG的调度频率、ID等业务信息。

正确解决方案

要获取MWAA中DAG的调度频率,有两种可靠方式:

方式1:通过MWAA ExecuteCommand调用Airflow CLI(推荐)

MWAA提供execute_command API,允许你在托管的Airflow环境中执行Airflow CLI命令,以此获取DAG元数据。以下是修改后的代码:

import boto3
import base64
import json
from botocore.config import Config

sso_profile = 'admin'
environment_name = 'dev-airflow'
config = Config(region_name='us-west-2')

session = boto3.Session(profile_name=sso_profile)
mwaa_client = session.client('mwaa', config=config)

try:
    # 执行Airflow CLI命令,获取所有DAG的JSON格式基础信息
    response = mwaa_client.execute_command(
        Name=environment_name,
        Command='airflow dags list --output json',
        ExecutionRoleArn='<你的MWAA执行角色ARN>'  # 替换为实际角色ARN
    )

    # 解码base64格式的输出内容
    stdout_data = base64.b64decode(response['Stdout']).decode('utf-8')
    if not stdout_data:
        print("No DAGs found in the MWAA environment")
    else:
        dags = json.loads(stdout_data)
        for dag in dags:
            # 执行命令获取单个DAG的详细信息(包含调度频率)
            dag_detail_response = mwaa_client.execute_command(
                Name=environment_name,
                Command=f'airflow dags show {dag["dag_id"]} --output json',
                ExecutionRoleArn='<你的MWAA执行角色ARN>'
            )
            dag_detail = json.loads(base64.b64decode(dag_detail_response['Stdout']).decode('utf-8'))
            print(f"DAG '{dag['dag_id']}' has a schedule interval of '{dag_detail['schedule_interval']}'")

except Exception as e:
    print(f"An error occurred: {e}")
注意事项
  • 确保你的执行角色拥有mwaa:ExecuteCommand权限,且该角色已关联到目标MWAA环境
  • 命令执行后的Stdout和Stderr均为base64编码,必须解码后才能解析使用
  • 若仅需调度频率,也可使用airflow dags list -sd命令,但dags show返回的信息更完整

方式2:直接读取S3中的DAG文件

如果你的DAG文件存储在S3桶中,可直接通过boto3读取文件内容,解析Python代码里的schedule_interval参数。示例代码:

import boto3
import ast

s3_client = boto3.client('s3')
bucket_name = '<你的DAG存储桶>'
dag_prefix = 'dags/'  # DAG文件在桶中的存储前缀

response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=dag_prefix)
for obj in response.get('Contents', []):
    if obj['Key'].endswith('.py'):
        dag_file = s3_client.get_object(Bucket=bucket_name, Key=obj['Key'])
        content = dag_file['Body'].read().decode('utf-8')
        # 简单解析静态定义的schedule_interval
        for line in content.split('\n'):
            if 'schedule_interval' in line and '#' not in line.split('=')[0]:
                schedule = line.split('=')[1].strip()
                # 尝试转换为Python原生对象
                try:
                    schedule = ast.literal_eval(schedule)
                except:
                    pass
                print(f"DAG file {obj['Key']} has schedule interval: {schedule}")
注意事项
  • 该方法仅适用于**静态定义schedule_interval**的DAG,若DAG是动态生成(比如从数据库读取调度规则),则无法正确获取
  • 需要确保执行角色拥有目标S3桶的读取权限

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 08:57:43