在Azure Data Factory中过滤Notebook Activity返回的空数组JSON对象
问题描述
我在Databricks Notebook中调用API获取输出,最终在Azure Data Factory(ADF)的Notebook Activity中得到JSON对象。现需过滤掉其中sub_indicator字段为空数组的整个JSON对象,仅保留该字段非空的对象。
输入JSON(来自Notebook Activity)
[ { "day": 60.0, "server": "xxxx", "database": "ddddd", "table": "tablename", "asset_id": "23232323", "indicate": ["value1"], "sub_indicator": ["sub"] }, {"day": 999.0, "server": "sadsadsad", "database": "dbbb", "table": "tablename2", "asset_id": "xxxxxx1", "indicate": ["value2"], "sub_indicator": ["sub2"] }, {"day": 30.0, "server": "server3", "database": "db3", "table": "tablename3", "asset_id": "xxxxxxx", "indicate": ["value3"], "sub_indicator": ["sub3"]}, {"day": 75.0, "server": "ser", "database": "db", "table": "tablename", "asset_id": "asdasd-adasdsa", "indicate": ["val1", "val2"], "sub_indicator": ["sub1", "sub2"]}, {"day": 50.0, "server": "serrr", "database": "dbb", "table": "tablename4", "asset_id": "yyyyyyyy", "indicate": ["value4"], "sub_indicator": ["sub1", "sub2"]}, {"day": 100.0, "server": "ser", "database": "IRF_Everest", "table": "tablename5", "asset_id": "adsadasdadasdasdasd", "indicate": ["sub1", "sub2"], "sub_indicator": ["val1"]}, {"day": 60.0, "server": "server3", "database": "db1", "table": "tablename7", "asset_id": "3312312321fsdasfasf", "indicate": ["val1"], "sub_indicator": []}, {"day": 50.0, "server": "serrrrr", "database": "db11", "table": "tablename8", "asset_id": "6ac9aea1-sdsdsdsadasdsadsad", "indicate": ["val"], "sub_indicator": []}, {"day": 60.0, "server": "serrr", "database": "db22", "table": "tablename10", "asset_id": "98e3dff0-adsadsadasd", "indicate": ["key"], "sub_indicator": ["sub_key"] } ]
待过滤对象示例
{ "day": 60.0, "server": "server3", "database": "db1", "table": "tablename7", "asset_id": "3312312321fsdasfasf", "indicate": ["val1"], "sub_indicator": [] }
期望输出JSON
[ { "day": 60.0, "server": "xxxx", "database": "ddddd", "table": "tablename", "asset_id": "23232323", "indicate": ["value1"], "sub_indicator": ["sub"] }, {"day": 999.0, "server": "sadsadsad", "database": "dbbb", "table": "tablename2", "asset_id": "xxxxxx1", "indicate": ["value2"], "sub_indicator": ["sub2"] }, {"day": 30.0, "server": "server3", "database": "db3", "table": "tablename3", "asset_id": "xxxxxxx", "indicate": ["value3"], "sub_indicator": ["sub3"]}, {"day": 75.0, "server": "ser", "database": "db", "table": "tablename", "asset_id": "asdasd-adasdsa", "indicate": ["val1", "val2"], "sub_indicator": ["sub1", "sub2"]}, {"day": 50.0, "server": "serrr", "database": "dbb", "table": "tablename4", "asset_id": "yyyyyyyy", "indicate": ["value4"], "sub_indicator": ["sub1", "sub2"]}, {"day": 100.0, "server": "ser", "database": "IRF_Everest", "table": "tablename5", "asset_id": "adsadasdadasdasdasd", "indicate": ["sub1", "sub2"], "sub_indicator": ["val1"] }, {"day": 60.0, "server": "serrr", "database": "db22", "table": "tablename10", "asset_id": "98e3dff0-adsadsadasd", "indicate": ["key"], "sub_indicator": ["sub_key"] } ]
解决方案
方法一:在Databricks Notebook中直接过滤(推荐)
既然数据从Databricks Notebook输出,直接在Notebook内完成过滤是最高效的方式,无需在ADF中额外处理。以下是Python实现代码:
# 假设API返回的原始数据存储在api_response变量中 api_response = [ # 原始JSON数组数据 ] # 过滤逻辑:保留sub_indicator非空的对象 filtered_data = [item for item in api_response if len(item.get('sub_indicator', [])) > 0] # 将过滤后的数据输出,供ADF的Notebook Activity接收 import json print(json.dumps(filtered_data))
说明:
- 用列表推导式遍历每个对象,检查
sub_indicator数组长度是否大于0; item.get('sub_indicator', [])确保对象无该字段时不报错,默认按空数组处理;- 输出的过滤后数据会直接被ADF的Notebook Activity获取。
方法二:在ADF中使用Filter Activity处理
若无法修改Databricks Notebook,可在ADF管道中添加Filter Activity处理Notebook输出:
- 将Notebook Activity的输出作为Filter Activity的输入;
- 配置Filter Activity的Condition为以下表达式:
该表达式会检查每个JSON对象的@greater(length(item().sub_indicator), 0)sub_indicator数组长度是否大于0; - 运行管道后,Filter Activity的输出即为过滤后的JSON数组。
说明:
item()代表数组中的单个元素;length()获取数组长度,greater()判断长度是否大于0,满足条件的对象会被保留。
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

