如何获取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
相关产品推荐
相关产品推荐

