如何在KSQL中创建可保留最新非空传感器值的表
问题描述
当前有接入多路传感器数据的流,传感器仅在自身状态更新时发送一次状态码。该值为一次性上报值,发送完成后传感器值会变为空,直至下一次状态变更。需求为让表中的最新有效值自动填充后续空值区间,直到新的有效值送达。
当前使用如下语句创建表:
CREATE TABLE LRS WITH (KAFKA_TOPIC='lrs', KEY_FORMAT='DELIMITED', PARTITIONS=6, REPLICAS=3) AS SELECT Device, LATEST_BY_OFFSET(CAST(Sensor1 AS DOUBLE)), LATEST_BY_OFFSET(CAST(Sensor2 AS DOUBLE)) FROM RELEVANT_VALUES RELEVANT_VALUES WINDOW TUMBLING ( SIZE 10 SECONDS ) GROUP BY Device
当前实现的实际输出效果如下:
Device | Sensor1 | Sensor2 | Timestamp 1 | null | null | 05:00am 1 | 3 | 2 | 05:01am 1 | null | null | 05:02am 1 | null | null | 05:03am 1 | 2 | 1 | 05:04am 1 | null | null | 05:05am
期望实现值持续更新填充空窗口的效果,具体如下:
Device | Sensor1 | Sensor2 | window 1 | null | null | 05:00-01 1 | 3 | 2 | 05:01-02 1 | 3 | 2 | 05:02-03 1 | 3 | 2 | 05:03-04 1 | 2 | 1 | 05:04-05 1 | 2 | 1 | 05:05-06
核心需求为创建一个始终展示最新上报非空值的KSQL Table,确认KSQL是否支持实现该需求。
解答
KSQL完全支持实现该需求,当前写法存在两个核心问题导致效果不符合预期:
- 滚动窗口(TUMBLING WINDOW)会按时间边界隔离计算状态,每个窗口独立做聚合,无法跨窗口继承上一个窗口的有效值
- 默认调用
LATEST_BY_OFFSET时会将null值作为有效值纳入计算,传感器上报的空值会直接覆盖之前留存的正常数值
可以根据实际输出粒度需求选择对应调整方案:
方案1:不需要严格按固定10秒窗口切分输出,仅需每个设备始终返回最新非空值
直接去掉窗口子句,同时开启LATEST_BY_OFFSET的忽略空值配置即可,修改后的SQL如下:
CREATE TABLE LRS WITH (KAFKA_TOPIC='lrs', KEY_FORMAT='DELIMITED', PARTITIONS=6, REPLICAS=3) AS SELECT Device, LATEST_BY_OFFSET(CAST(Sensor1 AS DOUBLE), true) AS Sensor1, LATEST_BY_OFFSET(CAST(Sensor2 AS DOUBLE), true) AS Sensor2 FROM RELEVANT_VALUES GROUP BY Device;
LATEST_BY_OFFSET的第二个参数设为true时,函数会自动忽略null值,只保留最新的非空上报结果,不会被后续的空值上报覆盖。
如果你使用的ksqldb版本低于0.17.0(该版本才新增上述忽略空值参数),可以先在聚合前过滤掉Sensor1、Sensor2全为null的消息,再做聚合计算,效果一致。
方案2:必须保留10秒滚动窗口的输出粒度
分两层处理实现:
- 先按照方案1的写法创建非窗口的设备状态表,持续维护每个设备的最新非空传感器值
- 再创建10秒周期的时钟驱动流(可通过定时任务往专用topic写入触发消息,或使用ksqldb内置的时间生成函数构造),将时钟流和设备状态表做关联,每个10秒窗口触发时直接读取当前状态表中存储的最新值输出,即可实现每个空窗口都自动填充上一个有效值的效果。
内容的提问来源于stack exchange,提问作者Tiscit
相关产品推荐
相关产品推荐

