You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从Sensor任务获取Xcom值并用于Dataflow任务参数

解决方法

步骤1:导入XComArg

首先确保导入Airflow的XComArg(适用于Airflow 2.x及以上版本):

from airflow.utils.xcom import XComArg

步骤2:从Sensor任务提取目标字符串

假设你的Sensor任务task_id为sensor_task,根据XCom的实际结构提取'df-00...'字符串:

  • 如果Sensor返回的XCom是字典,目标字符串在指定键下(比如file_name,需根据你的XCom结构调整键名):
    target_file = XComArg('sensor_task')['file_name']
    
  • 如果Sensor直接返回的就是'df-00...'字符串:
    target_file = XComArg('sensor_task')
    

步骤3:传入Dataflow任务参数

修改DataflowTemplatedJobStartOperator的配置,将提取到的字符串赋值给inputFilePattern:

dataflow_gcs_to_bq = DataflowTemplatedJobStartOperator(
    task_id='dataflow_gcs_to_bq',
    template='gs://your-template-path',
    job_name='your-job-name',
    location='your-region',
    project_id='your-project-id',
    parameters={
        'javascriptTextTransformFunctionName': 'your-transform-function',
        'JSONPath': 'gs://your-json-path',
        'javascriptTextTransformGcsPath': 'gs://your-js-path',
        'inputFilePattern': target_file,
        'outputTable': 'your-project:your-dataset.your-table',
        'bigQueryLoadingTemporaryDirectory': 'gs://your-temp-bucket',
    },        
)

注意事项

  • 替换代码中所有占位符(如sensor_task、your-template-path等)为实际业务值
  • 可通过Airflow UI的XCom页面确认Sensor返回的具体结构,调整提取逻辑

内容的提问来源于stack exchange,提问作者Jopsiton

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 21:32:54