Airflow如何引用上游任务输出列表的指定索引值?
解决Airflow动态映射中提取上游列表第一个元素的问题
直接使用upstream_task.output[0]会报错,因为output是代表上游任务完整输出的XComArg对象,下标访问会试图获取整个输出集合的第一个元素,而非每个子列表的首个元素。要实现动态映射时提取每个子列表的第一个元素,可通过以下方法解决:
方法1:用map处理上游输出
通过map函数遍历上游返回的每个子列表,提取第一个元素后再传给expand,适配现有代码逻辑:
修改后的start_instance任务代码如下:
start_instance = EC2StartInstanceOperator.partial( task_id='start-instance', region_name='eu-central-1', ).expand( instance_id=run_get_files.output.map(lambda item: item[0]), )
方法2:修改上游任务返回值(可选)
如果上游任务的输出仅需给start_instance提供第一个元素,可直接修改get_files函数,只返回每个子列表的首个元素:
def get_files(): # 仅返回每个子列表的第一个元素 return ['23', '49']
这种情况下直接用run_get_files.output传给expand即可,但如果上游输出还要给其他任务(比如run_create_ec2_instance)复用,不建议采用此方法。
原代码报错原因
run_get_files.output返回的是包含嵌套列表的XComArg集合,run_get_files.output[0]会试图获取整个集合的第一个元素(即['23', 'abcd']),而非遍历每个子列表提取首个值,不符合动态映射的需求,因此触发错误。
内容的提问来源于stack exchange,提问作者romanzdk
相关产品推荐
相关产品推荐

