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

