Airflow中Dynamic Task Mapping配置SimpleHttpOperator动态URI报错解决
Airflow动态任务映射配置动态URI报错解决
问题背景
我在Airflow中尝试通过Dynamic Task Mapping实现动态URI,DAG执行流程为 fetch_types >> request_tasks(request_tasks是动态映射任务)。其中fetch_types任务通过PostgreSQL Operator从数据库拉取类型数据,编写代码后运行报错。
错误代码
request_tasks = SimpleHttpOperator.partial( task_id='request_sync_task', http_conn_id='http_conn', endpoint='/api/v1/{{ type }}', method='GET', headers={"Content-Type": "application/json"}, response_check=lambda response: response.status_code == 200, dag=dag, ).expand( type="{{ task_instance.xcom_pull(task_ids='fetch_types', key='return_value') | map(attribute='1') }}" )
报错信息
SimpleHttpOperator.expand() got an unexpected keyword argument 'type'
不清楚如何配置URI路径以正确使用type值。
解决方法
type并非SimpleHttpOperator支持的可动态映射参数,将动态逻辑迁移到endpoint参数上,通过expand直接生成完整的URI路径即可解决问题。
修正后的示例代码(两种实现方式):
方式一:直接生成完整endpoint路径
request_tasks = SimpleHttpOperator.partial( task_id='request_sync_task', http_conn_id='http_conn', method='GET', headers={"Content-Type": "application/json"}, response_check=lambda response: response.status_code == 200, dag=dag, ).expand( endpoint="{{ task_instance.xcom_pull(task_ids='fetch_types', key='return_value') | map(attribute='1') | map('regex_replace', '^', '/api/v1/') | list }}" )
方式二:通过Jinja模板拼接路径
request_tasks = SimpleHttpOperator.partial( task_id='request_sync_task', http_conn_id='http_conn', endpoint='/api/v1/{{ item }}', method='GET', headers={"Content-Type": "application/json"}, response_check=lambda response: response.status_code == 200, dag=dag, ).expand( item="{{ task_instance.xcom_pull(task_ids='fetch_types', key='return_value') | map(attribute='1') }}" )
内容的提问来源于stack exchange,提问作者Euphemia
相关产品推荐
相关产品推荐

