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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 03:36:57