使用KSQL时为何表保留旧ROWTIME数据却丢弃ROWTIME更新的新数据?
问题根因
- ksqlDB 表默认按照Kafka消息的分区偏移量顺序更新主键对应的状态,与你自定义的
ROWTIME(即DateTime字段)无关。简单来说,后写入Kafka(偏移量更大)的消息无论时间戳大小,都会覆盖同主键之前的状态,你反复灌入测试数据时,如果旧时间戳的消息后写入Topic,自然会把之前的新时间戳状态覆盖,最终得到旧值。 - 你使用的
EMIT CHANGES是推送查询语法,会输出所有的状态变更记录,而非仅返回当前最新状态,所以你会看到同主键的多条历史变更输出,这并非是表内存储了多条同主键数据,只是查询将变更过程全量返回了。
解决方案
1. 开启时间戳排序更新(推荐)
如果你希望表仅保留同主键下时间戳更大的最新值,忽略乱序到达的旧时间戳消息,只需要在建表的WITH参数中新增TIMESTAMP_ORDERING = 'ENABLE'配置即可,修改后的建表语句如下:
CREATE TABLE vehicle_updates ( Latitude DOUBLE, Longitude DOUBLE, DateTime BIGINT, Registration STRING PRIMARY KEY ) WITH ( KAFKA_TOPIC = 'vehicle-update-log', VALUE_FORMAT = 'JSON_SR', TIMESTAMP = 'DateTime', TIMESTAMP_ORDERING = 'ENABLE' );
该配置开启后,ksqlDB处理消息时会先对比新消息的ROWTIME与当前主键已存储状态的ROWTIME,只有新消息时间戳更大时才会更新状态,自动忽略乱序的旧数据。
2. 使用拉取查询直接获取最新状态
如果你不需要跟踪变更过程,只需要查询某一时刻主键对应的最新值,不需要加EMIT CHANGES,直接使用拉取查询即可,语法如下:
SELECT registration, ROWTIME, TIMESTAMPTOSTRING(ROWTIME, 'yyyy-MM-dd HH:mm:ss.SSS', 'Africa/Johannesburg') AS rowtime_formatted FROM vehicle_updates WHERE registration = 'BT66MVE';
拉取查询只会返回当前表中该主键对应的唯一一条最新状态,不会返回历史变更记录。
3. 低版本兼容方案
如果你使用的是0.18版本之前不支持TIMESTAMP_ORDERING的ksqlDB,可以通过流处理分组聚合的方式实现取最新值的逻辑,示例如下:
CREATE STREAM vehicle_update_stream ( Latitude DOUBLE, Longitude DOUBLE, DateTime BIGINT, Registration STRING ) WITH ( KAFKA_TOPIC = 'vehicle-update-log', VALUE_FORMAT = 'JSON_SR', TIMESTAMP = 'DateTime' ); CREATE TABLE latest_vehicle_updates AS SELECT Registration, LATEST_BY_OFFSET(Latitude) Latitude, LATEST_BY_OFFSET(Longitude) Longitude, MAX(DateTime) DateTime FROM vehicle_update_stream GROUP BY Registration EMIT CHANGES;
内容的提问来源于stack exchange,提问作者Pieter Breed
相关产品推荐
相关产品推荐

