如何提取Pipeline中使用的数据集名称以确定输出数据安全分级?
自动获取数据流依赖数据集并设置最高安全分级
一、用Pipeline+GetMetadata解决数据集名称获取问题
你之前没拿到数据集列表,大概率是GetMetadata的配置没到位,按以下步骤调整:
- 构建Pipeline,先添加目标数据流活动(比如
dataFlow1) - 新增GetMetadata活动,设置为数据流执行完成后触发
- 在GetMetadata的配置页:
- 数据源类型选
Data Flow,选中你要分析的数据流dataFlow1 - 字段列表里勾选
datasets,这个字段会返回数据流中所有引用的输入数据集名称数组
- 数据源类型选
- 运行Pipeline后,查看GetMetadata的输出,就能拿到完整的数据集名称列表
二、查询分级并计算最高值
拿到数据集列表后,需要关联Classification表获取分级并取最大值:
- 新增Set Variable活动,创建字符串变量
datasetListStr,用表达式把数据集数组转成SQL IN子句需要的格式:
比如数组@concat("'", join(activity('GetDataFlowDatasets').output.datasets, "','"), "'")["datasetA","datasetB"]会转成'datasetA','datasetB' - 新增Lookup活动,查询Classification表,SQL语句写:
SELECT security_level FROM Classification WHERE dataset_name IN (@variables('datasetListStr')) - 再新增Set Variable活动,创建整数变量
maxSecurityLevel,用表达式取Lookup结果中的最大值:@max(activity('LookupSecurityLevels').output.value, item().security_level)
三、写入新数据集的分级记录
最后把新生成的数据集名称和最高分级插入Classification表:
- 可以用Copy Data活动,将包含新数据集名称和
maxSecurityLevel的临时数据写入Classification表 - 或者用Stored Procedure活动,调用预先写好的存储过程,传入新数据集名称和最高分级参数
四、Notebook方案的修正
如果想用Notebook实现,核心是通过ADF API获取数据流元数据,以下是PySpark示例代码:
import requests import json # 替换为你的Azure资源信息 subscription_id = "你的订阅ID" adf_name = "你的数据工厂名称" resource_group = "你的资源组名称" data_flow_name = "dataFlow1" # 获取访问令牌(可通过Azure CLI执行az account get-access-token --resource https://management.azure.com获取) access_token = "你的访问令牌" # 调用ADF API获取数据流元数据 url = f"https://management.azure.com/subscriptions/{subscription_id}/resourceGroups/{resource_group}/providers/Microsoft.DataFactory/factories/{adf_name}/dataflows/{data_flow_name}?api-version=2018-06-01" headers = {"Authorization": f"Bearer {access_token}", "Content-Type": "application/json"} response = requests.get(url, headers=headers) data_flow_meta = response.json() # 提取所有输入数据集名称 input_datasets = [] for source in data_flow_meta['properties']['sources']: if 'dataset' in source: input_datasets.append(source['dataset']['referenceName']) input_datasets = list(set(input_datasets)) # 去重 # 查询Classification表获取分级 df_classification = spark.sql("SELECT dataset_name, security_level FROM Classification WHERE dataset_name IN ('" + "','".join(input_datasets) + "')") # 计算最高分级 max_level = df_classification.agg({"security_level": "max"}).collect()[0][0] # 写入新记录到Classification表 new_dataset_name = "你的输出数据集名称" spark.sql(f"INSERT INTO Classification (dataset_name, security_level) VALUES ('{new_dataset_name}', {max_level})")
注意:Notebook需要有访问ADF API和Classification表的权限,否则会运行失败。
内容的提问来源于stack exchange,提问作者Christopher Lindeman
相关产品推荐
相关产品推荐

