如何在Apache NiFi中编写Jolt Spec实现JSON数组聚合
Apache NiFi中基于deviceid分组并统计记录数的Jolt实现
输入示例
[ { "altitude": 19.1, "analog1": 0.016, "analog2": 0.004, "batchprocessgroupid": 0, "batterylevel": 1.598, "bleid": "", "command": 3, "devicedatetime": 1679325571000, "deviceid": "863071015139949", "digital1": 0, "digital2": 0, "gpsdirection": 160.1, "gpsfix": 1, "hdop": 1, "ibutton": "47000019DBC1D001", "ignitionstatus": 1, "ioeventid": 7, "iostring": "5 -> 1, 78 -> 2003, 29 -> 28764, 28 -> 1, 70 -> 0000000000000000, 137 -> 0, 65 -> 164295509, 173 -> 1, 134 -> 0, 2 -> 0, 32 -> 38, 34 -> 47000019DBC1D001, 22 -> 16, 27 -> 14, 71 -> 0000000000000000, 49 -> 0, 130 -> 2, 135 -> 0, 3 -> 0, 150 -> 42402, 139 -> 60, 23 -> 4, 30 -> 1598, 136 -> 0, 79 -> 2003, 131 -> 0", "isprimary": "1", "latitude": 24.99225, "listenerdatetime": 1679311155585, "locationbit": 1, "locationen": "Dubai, Mena Jabal Ali", "longitude": 55.0920783, "mainpower": 28.764, "odometer": 164295.509, "panic": 0, "providertenantuids": "", "reason": 0, "recordstatus": 1, "satellites": 11, "speed": 0, "temperature1": 0, "temperature2": 0, "tenantgroupuid": "4", "tenantuid": "2", "uniqueid": 19532 }, { "altitude": 16.3, "analog1": 0.016, "analog2": 0.005, "batchprocessgroupid": 0, "batterylevel": 2.645, "bleid": "", "command": 3, "devicedatetime": 1679325626000, "deviceid": "863071015139949", "digital1": 0, "digital2": 0, "gpsdirection": 307.9, "gpsfix": 1, "hdop": 1, "ibutton": "47000019DBC1D001", "ignitionstatus": 1, "ioeventid": 9, "iostring": "5 -> 1, 78 -> 2003, 29 -> 28793, 28 -> 1, 70 -> 0000000000000000, 137 -> 0, 65 -> 164295532, 173 -> 1, 134 -> 0, 2 -> 0, 32 -> 38, 34 -> 47000019DBC1D001, 22 -> 16, 27 -> 20, 71 -> 0000000000000000, 49 -> 255, 130 -> 5, 135 -> 0, 3 -> 0, 150 -> 42402, 139 -> 55, 23 -> 5, 30 -> 2645, 136 -> 0, 79 -> 2003, 131 -> 0", "isprimary": "1", "latitude": 24.9923233, "listenerdatetime": 1679311210701, "locationbit": 1, "locationen": "Dubai, Mena Jabal Ali", "longitude": 55.0921766, "mainpower": 28.793, "odometer": 164295.532, "panic": 0, "providertenantuids": "", "reason": 0, "recordstatus": 1, "satellites": 11, "speed": 5, "temperature1": 0, "temperature2": 0, "tenantgroupuid": "4", "tenantuid": "2", "uniqueid": 31628 } ]
期望输出
[ { "deviceid": "863071015139949", "count": 2, "devicedata": [ { "altitude": 19.1, "analog1": 0.016, "analog2": 0.004, "batchprocessgroupid": 0, "batterylevel": 1.598, "bleid": "", "command": 3, "devicedatetime": 1679325571000, "deviceid": "863071015139949", "digital1": 0, "digital2": 0, "gpsdirection": 160.1, "gpsfix": 1, "hdop": 1, "ibutton": "47000019DBC1D001", "ignitionstatus": 1, "ioeventid": 7, "iostring": "5 -> 1, 78 -> 2003, 29 -> 28764, 28 -> 1, 70 -> 0000000000000000, 137 -> 0, 65 -> 164295509, 173 -> 1, 134 -> 0, 2 -> 0, 32 -> 38, 34 -> 47000019DBC1D001, 22 -> 16, 27 -> 14, 71 -> 0000000000000000, 49 -> 0, 130 -> 2, 135 -> 0, 3 -> 0, 150 -> 42402, 139 -> 60, 23 -> 4, 30 -> 1598, 136 -> 0, 79 -> 2003, 131 -> 0", "isprimary": "1", "latitude": 24.99225, "listenerdatetime": 1679311155585, "locationbit": 1, "locationen": "Dubai, Mena Jabal Ali", "longitude": 55.0920783, "mainpower": 28.764, "odometer": 164295.509, "panic": 0, "providertenantuids": "", "reason": 0, "recordstatus": 1, "satellites": 11, "speed": 0, "temperature1": 0, "temperature2": 0, "tenantgroupuid": "4", "tenantuid": "2", "uniqueid": 19532 }, { "altitude": 16.3, "analog1": 0.016, "analog2": 0.005, "batchprocessgroupid": 0, "batterylevel": 2.645, "bleid": "", "command": 3, "devicedatetime": 1679325626000, "deviceid": "863071015139949", "digital1": 0, "digital2": 0, "gpsdirection": 307.9, "gpsfix": 1, "hdop": 1, "ibutton": "47000019DBC1D001", "ignitionstatus": 1, "ioeventid": 9, "iostring": "5 -> 1, 78 -> 2003, 29 -> 28793, 28 -> 1, 70 -> 0000000000000000, 137 -> 0, 65 -> 164295532, 173 -> 1, 134 -> 0, 2 -> 0, 32 -> 38, 34 -> 47000019DBC1D001, 22 -> 16, 27 -> 20, 71 -> 0000000000000000, 49 -> 255, 130 -> 5, 135 -> 0, 3 -> 0, 150 -> 42402, 139 -> 55, 23 -> 5, 30 -> 2645, 136 -> 0, 79 -> 2003, 131 -> 0", "isprimary": "1", "latitude": 24.9923233, "listenerdatetime": 1679311210701, "locationbit": 1, "locationen": "Dubai, Mena Jabal Ali", "longitude": 55.0921766, "mainpower": 28.793, "odometer": 164295.532, "panic": 0, "providertenantuids": "", "reason": 0, "recordstatus": 1, "satellites": 11, "speed": 5, "temperature1": 0, "temperature2": 0, "tenantgroupuid": "4", "tenantuid": "2", "uniqueid": 31628 } ] } ]
Jolt转换规则
使用以下Jolt规范可实现按deviceid分组并统计记录数量的需求:
[ { "operation": "shift", "spec": { "*": { "@": "&1.devicedata[]", "deviceid": "&1.deviceid" } } }, { "operation": "shift", "spec": { "*": { "deviceid": "@(1,deviceid).deviceid", "devicedata": "@(1,deviceid).devicedata" } } }, { "operation": "shift", "spec": { "*": { "@": "[]" } } }, { "operation": "modify-overwrite-beta", "spec": { "*": { "count": "=size(@(1,devicedata))" } } } ]
规则说明
- 第一步(shift操作):将输入数组中的每个对象,完整映射到临时索引键下的
devicedata数组,同时单独提取deviceid字段。 - 第二步(shift操作):以
deviceid作为分组键,把同一设备的所有数据聚合到对应的键下,完成分组。 - 第三步(shift操作):将分组后的键值对转换为数组结构,匹配期望输出的外层数组格式。
- 第四步(modify-overwrite-beta操作):调用
size函数计算devicedata数组的长度,生成count字段记录该设备的记录总数。
内容的提问来源于stack exchange,提问作者naga babu
相关产品推荐
相关产品推荐

