如何从ksqlDB的QUERYABLE_TESTTABLE聚合表中删除记录?
问题场景
我们有如下Kafka与ksqlDB环境:
- Kafka主题
customer_events,消息值示例:
{ "CUSTOMERID": "198fa518-1031-4fe8-8abd-ca29bd120259" }
- 在ksqlDB中基于该主题创建持久流:
CREATE STREAM TEST_STREAM (SESSIONID STRING KEY, CUSTOMERID STRING) WITH (KAFKA_TOPIC='customer_events', KEY_FORMAT='KAFKA', PARTITIONS=1, VALUE_FORMAT='JSON');
- 基于流创建聚合表,按
SESSIONID聚合客户ID列表:
CREATE TABLE QUERYABLE_TESTTABLE AS SELECT SRC.SESSIONID SESSIONID, COLLECT_LIST(SRC.CUSTOMERID) CUSTOMERS FROM TEST_STREAM SRC GROUP BY SRC.SESSIONID EMIT CHANGES;
- 拉取查询表的结果符合预期,但无法删除表中指定条目,尝试向
customer_events或表的底层主题写入墓碑消息均无效。
解决方案
方法一:使用ksqlDB DELETE语句(推荐,适用于0.18及以上版本)
ksqlDB 0.18版本起支持直接对物化表执行删除操作,语句会自动清理状态存储并发送墓碑消息:
DELETE FROM QUERYABLE_TESTTABLE WHERE SESSIONID = '目标SessionId值';
执行后,该SessionId对应的条目会从表中移除,后续查询不再返回。
方法二:配置状态存储TTL自动清理
如果希望自动清理长时间无更新的条目,可在创建表时设置状态存储的过期时间(单位毫秒):
CREATE TABLE QUERYABLE_TESTTABLE AS SELECT SRC.SESSIONID SESSIONID, COLLECT_LIST(SRC.CUSTOMERID) CUSTOMERS FROM TEST_STREAM SRC GROUP BY SRC.SESSIONID EMIT CHANGES WITH (STATE_STORE_TTL='86400000'); -- 示例:1天过期
当某个SessionId超过设置的时间没有新消息更新时,ksqlDB会自动清理该条目,底层主题也会生成对应的墓碑消息。
注意事项
直接向原流或表的底层主题写入墓碑消息对全局聚合表无效,因为全局聚合的状态默认永久保留,仅靠墓碑消息无法触发聚合状态的清理,必须使用上述两种方法。
内容的提问来源于stack exchange,提问作者Pavel Cermak
相关产品推荐
相关产品推荐

