Airflow中WasbBlobSensor如何推送XCom传递blob名称以复用SparkSubmitOperator
WasbBlobSensor 推送 XCom 及工作流适配方案
方案1:自定义扩展 WasbBlobSensor 直接推送 blob 名到 XCom
原生 WasbBlobSensor 本身支持 do_xcom_push 参数(默认值为 True),但默认只会将 poke 方法返回的布尔值(检测到文件为 True)推送到 XCom,所以只需要继承该传感器重写 poke 方法,将匹配到的 blob 名称作为返回值即可:
from airflow.providers.microsoft.azure.sensors.wasb import WasbBlobSensor from airflow.utils.context import Context class WasbBlobSensorWithXCom(WasbBlobSensor): def poke(self, context: Context) -> str | bool: # 调用父类原有检测逻辑 blob_exists = super().poke(context) if blob_exists: # 检测到文件存在时返回blob名称,会自动推送到XCom的return_value字段 return self.blob_name return False
替换你原有传感器的生成逻辑即可:
blob_sensors = [ WasbBlobSensorWithXCom( task_id=blob.split('.')[0] + "_sensor", wasb_conn_id="wasb_default", container_name=blob_container, blob_name=blob, poke_interval=60, timeout=120, do_xcom_push=True, dag=dag) for blob in blob_names]
之后在复用的 SparkSubmitOperator 中,直接通过模板语法拉取所有传感器的 XCom 值即可拿到触发成功的对应 blob 名:
spark_job = SparkSubmitOperator( task_id="spark_import_mysql", application="your_spark_job.py", # 模板语法拉取所有传感器的XCom,成功触发的传感器会返回blob名,其余为None application_args=["--blob-name", "{{ ti.xcom_pull(task_ids=[s.task_id for s in blob_sensors]) | select('!=', None) | first }}"], trigger_rule="one_success", dag=dag ) # 所有传感器都指向该Spark作业 for sensor in blob_sensors: sensor >> spark_job
方案2:不自定义传感器,新增过渡 PythonOperator 传递 blob 名
如果你不想修改传感器原有逻辑,可以给每个传感器绑定一个轻量 PythonOperator 负责推送对应 blob 名到 XCom:
from airflow.operators.python import PythonOperator def push_blob_name(blob_name, **context): return blob_name blob_sensors = [] push_tasks = [] for blob in blob_names: sensor = WasbBlobSensor( task_id=blob.split('.')[0] + "_sensor", wasb_conn_id="wasb_default", container_name=blob_container, blob_name=blob, poke_interval=60, timeout=120, dag=dag ) push_task = PythonOperator( task_id=f"push_{blob.split('.')[0]}_name", python_callable=push_blob_name, op_kwargs={"blob_name": blob}, do_xcom_push=True, dag=dag ) sensor >> push_task blob_sensors.append(sensor) push_tasks.append(push_task) # Spark作业拉取所有push任务的XCom即可 spark_job = SparkSubmitOperator( task_id="spark_import_mysql", application="your_spark_job.py", application_args=["--blob-name", "{{ ti.xcom_pull(task_ids=[t.task_id for t in push_tasks]) | select('!=', None) | first }}"], trigger_rule="one_success", dag=dag ) for push_task in push_tasks: push_task >> spark_job
目标表匹配逻辑
拿到 blob 名后,你可以在 Spark 代码中提前维护 blob 名和目标 MySQL 表的映射关系,直接通过传入的 blob 名匹配到对应表执行导入即可。
内容的提问来源于stack exchange,提问作者Minura Punchihewa
相关产品推荐
相关产品推荐

