Azure Data Factory实现:启动Pipeline A时停止Pipeline B直至其完成
实现思路
核心逻辑分为5个步骤:
- 检查Pipeline B是否处于运行状态
- 若B正在运行,立即终止所有运行实例
- 执行Pipeline A的全部业务逻辑
- 等待Pipeline A执行完成
- 启动Pipeline B
REST API调用细节
需用到三个ADF REST API,均通过Web Activity调用,且需配置**托管标识(Managed Identity)**认证(在Web Activity的Authentication选项中选择Managed Identity,授予ADF实例Data Factory Contributor权限):
1. 查询运行中的Pipeline B实例
- 请求方法:
GET - 请求URL(替换占位符为实际信息):
https://management.azure.com/subscriptions/{你的订阅ID}/resourceGroups/{你的资源组名}/providers/Microsoft.DataFactory/factories/{你的ADF工厂名}/pipelineruns?api-version=2018-06-01&$filter=pipelineName eq '{PipelineB名称}' and status eq 'InProgress'
- 返回结果中
value数组为空则说明B未在运行,反之则包含所有运行中实例的ID。
2. 终止指定Pipeline B运行实例
- 请求方法:
POST - 请求URL(替换占位符及
{runId}为查询到的实例ID):
https://management.azure.com/subscriptions/{你的订阅ID}/resourceGroups/{你的资源组名}/providers/Microsoft.DataFactory/factories/{你的ADF工厂名}/pipelineruns/{runId}/cancel?api-version=2018-06-01
- 若存在多个运行实例,需用ForEach Activity遍历
value数组逐个终止。
3. 启动Pipeline B
- 请求方法:
POST - 请求URL(替换占位符):
https://management.azure.com/subscriptions/{你的订阅ID}/resourceGroups/{你的资源组名}/providers/Microsoft.DataFactory/factories/{你的ADF工厂名}/pipelines/{PipelineB名称}/createRun?api-version=2018-06-01
- 需传递参数时,可在请求体中添加
parameters字段。
完善后的Pipeline配置示例
基于你提供的JSON补充关键配置:
{ "name": "p_Stop_go_Test_2", "properties": { "activities": [ { "name": "Lookup1", "type": "Lookup", "dependsOn": [], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "source": { "type": "AzureSqlSource", "queryTimeout": "02:00:00", "partitionOption": "None" }, "dataset": { "referenceName": "yyyyy", "type": "DatasetReference" } } }, { "name": "get pStopgoTest runs based on filters", "description": "get p_Stop_go_Test runs based on filters.", "type": "WebActivity", "dependsOn": [], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "url": "https://management.azure.com/subscriptions/{订阅ID}/resourceGroups/{资源组名}/providers/Microsoft.DataFactory/factories/{ADF工厂名}/pipelineruns?api-version=2018-06-01&$filter=pipelineName eq 'p_Stop_go_Test' and status eq 'InProgress'", "method": "GET", "authentication": { "type": "MSI", "resource": "https://management.azure.com/" } } }, { "name": "Check if response is empty or not", "description": "in the rest api response if the values array is empty then pStopgoTest is not in running state", "type": "IfCondition", "dependsOn": [ { "activity": "get pStopgoTest runs based on filters", "dependencyConditions": [ "Succeeded" ] } ], "userProperties": [], "typeProperties": { "expression": { "value": "@not(empty(activity('get pStopgoTest runs based on filters').output.value))", "type": "Expression" }, "ifFalseActivities": [ { "name": "执行Pipeline A核心逻辑", "type": "ExecutePipeline", "dependsOn": [], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "pipeline": { "referenceName": "PipelineA", "type": "PipelineReference" }, "waitOnCompletion": true } }, { "name": "启动Pipeline B", "type": "WebActivity", "dependsOn": [ { "activity": "执行Pipeline A核心逻辑", "dependencyConditions": [ "Succeeded" ] } ], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "url": "https://management.azure.com/subscriptions/{订阅ID}/resourceGroups/{资源组名}/providers/Microsoft.DataFactory/factories/{ADF工厂名}/pipelines/p_Stop_go_Test/createRun?api-version=2018-06-01", "method": "POST", "authentication": { "type": "MSI", "resource": "https://management.azure.com/" } } } ], "ifTrueActivities": [ { "name": "遍历运行中的Pipeline B实例", "type": "ForEach", "dependsOn": [], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "items": { "value": "@activity('get pStopgoTest runs based on filters').output.value", "type": "Expression" }, "activities": [ { "name": "Cancle Pipeline pStopgoTest", "type": "WebActivity", "dependsOn": [], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "url": { "value": "@concat('https://management.azure.com/subscriptions/{订阅ID}/resourceGroups/{资源组名}/providers/Microsoft.DataFactory/factories/{ADF工厂名}/pipelineruns/', item().runId, '/cancel?api-version=2018-06-01')", "type": "Expression" }, "method": "POST", "authentication": { "type": "MSI", "resource": "https://management.azure.com/" } } } ] } }, { "name": "执行Pipeline A核心逻辑", "type": "ExecutePipeline", "dependsOn": [ { "activity": "遍历运行中的Pipeline B实例", "dependencyConditions": [ "Succeeded" ] } ], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "pipeline": { "referenceName": "PipelineA", "type": "PipelineReference" }, "waitOnCompletion": true } }, { "name": "启动Pipeline B", "type": "WebActivity", "dependsOn": [ { "activity": "执行Pipeline A核心逻辑", "dependencyConditions": [ "Succeeded" ] } ], "policy": { "timeout": "0.12:00:00", "retry": 0, "retryIntervalInSeconds": 30, "secureOutput": false, "secureInput": false }, "userProperties": [], "typeProperties": { "url": "https://management.azure.com/subscriptions/{订阅ID}/resourceGroups/{资源组名}/providers/Microsoft.DataFactory/factories/{ADF工厂名}/pipelines/p_Stop_go_Test/createRun?api-version=2018-06-01", "method": "POST", "authentication": { "type": "MSI", "resource": "https://management.azure.com/" } } } ] } } ], "annotations": [], "lastPublishTime": "xxx" }, "type": "Microsoft.DataFactory/factories/pipelines" }
关键注意事项
- 确保ADF实例已启用托管标识,并授予Data Factory Contributor角色权限,否则REST API调用会触发权限不足错误。
ExecutePipeline活动的waitOnCompletion需设为true,保证Pipeline A完全执行完毕后再启动B。- 若Pipeline B存在多个并行运行实例,必须用ForEach遍历终止,避免遗漏。
内容的提问来源于stack exchange,提问作者cinnamonroll
相关产品推荐
相关产品推荐

