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

如何获取Cloud Data Fusion管道上次运行时间戳与状态并作为运行时参数传入

Cloud Data Fusion管道增量同步:获取上次运行时间戳与状态方案

方法1:通过Cloud Data Fusion REST API直接查询

Cloud Data Fusion原生没有内置宏获取上次运行信息,但可以通过官方REST API查询管道的历史运行记录,提取最新的一条数据:

  • 构造API请求,按结束时间倒序取第一条运行记录,拿到时间戳和状态:
curl -H "Authorization: Bearer $(gcloud auth print-access-token)" \
"https://REGION-datafusion.googleapis.com/v1/projects/PROJECT_ID/locations/REGION/instances/INSTANCE_ID/namespaces/NAMESPACE_ID/pipelines/PIPELINE_ID/runs?orderBy=endTime desc&pageSize=1"
  • 替换请求中的REGION、PROJECT_ID、INSTANCE_ID、NAMESPACE_ID、PIPELINE_ID为实际值;返回的JSON响应里,endTime字段是上次运行结束的时间戳,state字段是运行状态(如SUCCEEDED、FAILED)。
  • 可以在管道启动前写个脚本调用该API,把结果作为运行时参数传入管道启动命令:
gcloud data-fusion pipelines run PIPELINE_ID \
--instance INSTANCE_ID \
--location REGION \
--parameter last_run_time=2024-05-20T12:00:00Z \
--parameter last_run_state=SUCCEEDED

方法2:将运行信息持久化到外部存储,启动前读取

如果不想依赖API,可以把每次管道的运行信息存到外部存储,下次启动时读取:

  • 在管道末尾添加一个Shell执行组件,将本次运行的时间戳和状态写入Cloud Storage(或Firestore):
echo '{"last_run_time": "'$(date -u +%Y-%m-%dT%H:%M:%SZ)'", "state": "SUCCEEDED"}' | gsutil cp - gs://your-bucket/pipeline-metadata/last-run.json
  • 启动管道前,通过脚本读取存储的文件内容,提取参数后传入:
LAST_RUN=$(gsutil cat gs://your-bucket/pipeline-metadata/last-run.json)
LAST_RUN_TIME=$(echo $LAST_RUN | jq -r '.last_run_time')
LAST_RUN_STATE=$(echo $LAST_RUN | jq -r '.state')

gcloud data-fusion pipelines run PIPELINE_ID \
--instance INSTANCE_ID \
--location REGION \
--parameter last_run_time=$LAST_RUN_TIME \
--parameter last_run_state=$LAST_RUN_STATE

方法3:用Cloud Composer(Airflow)编排,利用Airflow变量存储

如果用Cloud Composer触发Data Fusion管道,可以借助Airflow的变量系统管理运行信息:

  • 在Airflow DAG中,执行完Data Fusion管道任务后,将运行时间和状态存入Airflow变量:
from airflow.models import Variable

# 假设task是Data Fusion管道任务
task_instance = context['task_instance']
Variable.set(
    "datafusion_pipeline_last_run",
    {
        "time": context['execution_date'].isoformat(),
        "state": task_instance.state
    }
)
  • 下次触发DAG时,读取变量并作为参数传入Data Fusion管道:
last_run_info = Variable.get("datafusion_pipeline_last_run", deserialize_json=True)

datafusion_run_task = DataFusionStartPipelineOperator(
    task_id="start_datafusion_pipeline",
    pipeline_name="PIPELINE_ID",
    instance_name="INSTANCE_ID",
    location="REGION",
    parameters={
        "last_run_time": last_run_info["time"],
        "last_run_state": last_run_info["state"]
    }
)

增量同步的实际使用

拿到last_run_time和last_run_state后,在管道的数据源组件(如JDBC、BigQuery)中添加过滤条件:

  • 仅当last_run_state为SUCCEEDED时,使用last_updated >= ${last_run_time}作为增量过滤条件;如果上次运行失败,可选择回退到更早的时间戳或执行全量同步,避免数据遗漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 16:02:28