Airflow TaskGroup中拉取列表型XCom失败问题排查
Airflow动态创建Dataproc集群的XCom使用问题
场景说明
通过Python Operator的可调用函数生成列表并推送到XCom,随后尝试拉取这些XCom数据来动态创建Dataproc集群,但过程中遇到Task ID相关错误。
推送XCom的Python代码
keys = [] values = [] def attribute_count_check(e_run_id,**context): job_run_id = int(e_run_id) da = "select count (distinct row_num) from dds_metadata.dds_temp_att_table where run_id ={}".format(job_run_id) cursor.execute(da) res = cursor.fetchall() view_res = [x for res in res for x in res] count_of_sql = view_res[0] print(count_of_sql) if count_of_sql < 1: print("deleting of cluster") return 'delete_cluster' else : print("triggering attr_check") num_attributes_per_task = num_attr # job_config diff = math.ceil(count_of_sql / num_attributes_per_task) instance = int(diff) n = num_attributes_per_task global values global keys for r in range(1, instance+1): keys.append(r) lower_ranges = (n*(r-1)) + 1 upper_range = (n*(r - 1)) + n b = (lower_ranges, upper_range) values.append(b) task_instance = context['task_instance'] task_instance.xcom_push(key="di_keys", value=keys) task_instance.xcom_push(key="di_values", value=values)
XCom存储的数据
已成功将di_keys(示例结构如[1,2,3])和di_values(示例结构如[(1,10),(11,20),(21,30)])推送到XCom中。
动态创建集群的初始代码
with TaskGroup('dataproc_create_cluster', prefix_group_id=False) as dataproc_create_clusters: for i in zip('{{ ti.xcom_pull(key="di_keys")}}','{{ ti.xcom_pull(key="di_values")}}'): dynmaic_create_cluster = DataprocCreateClusterOperator( task_id="create_cluster_{}".format(list(eval(str(i)))[0]), project_id='{0}'.format(PROJECT), cluster_config=CLUSTER_GENERATOR_CONFIG, region='{0}'.format(REGION), cluster_name="dataproc-cluster-{}-sit".format(str(i[0])), )
第一个错误信息
Broken DAG: [/opt/airflow/dags/Cluster_config.py] Traceback (most recent call last): File "/usr/local/lib/python3.6/site-packages/airflow/models/baseoperator.py", line 547, in __init__ validate_key(task_id) File "/usr/local/lib/python3.6/site-packages/airflow/utils/helpers.py", line 56, in validate_key "dots and underscores exclusively".format(k=k) airflow.exceptions.AirflowException: The key (create_cluster_{) has to be made of alphanumeric characters, dashes, dots and underscores exclusively
修改Task ID后的代码
为解决上述错误,修改了task_id的生成方式:
task_id="create_cluster_"+re.sub(r'\W+', '', str(list(eval(str(i)))[0])),
修改后出现的第二个错误
airflow.exceptions.DuplicateTaskIdFound: Task id 'create_cluster_' has already been added to the DAG
已尝试在相关配置中添加render_template_as_native_obj=True,但仍出现重复Task ID的错误,寻求解决方法。
内容的提问来源于stack exchange,提问作者djgcp
相关产品推荐
相关产品推荐

