Airflow使用AWSAthenaOperator执行动态查询列表报错如何解决
问题根因
Airflow的DAG结构是在DagBag加载阶段(也就是调度器扫描DAG文件的时候)就已经确定完成的,此时所有DAG定义代码都会被执行,但前序任务还没有实际运行,read_date_information_file.output只是一个XCom的占位引用,不是真实的返回列表。你直接用enumerate遍历这个占位符的时候,会触发Airflow尝试解析这个还不存在的XCom值,最终导致加载超时。
硬编码的列表可以正常运行,是因为DagBag加载阶段就能确定列表长度和内容,可以直接生成固定数量的Athena任务。
解决方案1:使用动态任务映射(推荐,Airflow 2.3+支持)
动态任务映射是Airflow官方推出的专门用来处理「根据前序任务的输出动态生成多个下游任务」的场景,不需要手动遍历列表,Airflow会在运行时根据前序任务返回的列表自动生成对应数量的任务。
修改后的代码示例:
def get_date_information(ti): import boto3 s3 = boto3.client('s3') data = s3.get_object(Bucket=output_bucket, Key=key) contents = data['Body'].read().decode("utf-8") print('Date information is: ', contents) events_list = contents.split(',') return events_list with DAG( dag_id='adserver_split_job_emr_job_dag', default_args={ 'owner': 'adserver_airflow', 'depends_on_past': False, 'email': ['airflow@example.com'], 'email_on_failure': False, 'email_on_retry': False, }, dagrun_timeout=timedelta(hours=2), start_date=datetime(2021, 9, 22, 9), schedule_interval='20 * * * *', # Airflow 2.4以下版本需要加这个参数开启动态任务映射 # render_template_as_native_obj=True ) as dag: read_date_information_file = PythonOperator( task_id="read_date_information_file", python_callable=get_date_information ) # partial定义所有任务通用的参数,expand指定需要动态遍历的参数 run_query = AWSAthenaOperator.partial( task_id='run_query', output_location=config.ATHENA_OUTPUT_LOCATION, database=config.ATHENA_DATABASE_NAME, aws_conn_id='aws_default' ).expand(query=read_date_information_file.output) read_date_information_file >> run_query
Airflow会自动给每个query生成独立的任务,任务ID会自动加上序号后缀,支持单个任务的重试、日志查看等能力。
解决方案2:老版本Airflow兼容方案
如果你的Airflow版本低于2.3,没法使用动态任务映射,可以把Athena查询的执行逻辑封装到一个Python函数里,在函数内部遍历查询列表挨个执行:
def run_athena_queries(ti): import boto3 from pyathena import connect query_list = ti.xcom_pull(task_ids='read_date_information_file') # 初始化Athena连接 conn = connect(s3_staging_dir=config.ATHENA_OUTPUT_LOCATION, region_name='你的AWS区域', aws_conn_id='aws_default') cursor = conn.cursor() for i, event in enumerate(query_list): print(f'执行第{i}个查询: {event}') cursor.execute(event) # 可自行添加等待查询完成、获取结果的逻辑 print(f'第{i}个查询执行完成') with DAG( dag_id='adserver_split_job_emr_job_dag', default_args={ 'owner': 'adserver_airflow', 'depends_on_past': False, 'email': ['airflow@example.com'], 'email_on_failure': False, 'email_on_retry': False, }, dagrun_timeout=timedelta(hours=2), start_date=datetime(2021, 9, 22, 9), schedule_interval='20 * * * *', ) as dag: read_date_information_file = PythonOperator( task_id="read_date_information_file", python_callable=get_date_information ) run_all_queries = PythonOperator( task_id='run_all_queries', python_callable=run_athena_queries ) read_date_information_file >> run_all_queries
注意:该方案的缺点是所有查询都在同一个任务里执行,没法单独重试单个查询,如果需要单个查询的任务级别控制,还是建议升级Airflow版本使用动态任务映射。
内容的提问来源于stack exchange,提问作者seou1
相关产品推荐
相关产品推荐

