如何获取Kafka Streams中拆分字段的最小值与最大值?
解决Kafka Streams/KSQL分组获取最小/最大值对应完整记录的问题
你的核心需求是按payload_1的前半部分(比如history::05000228023411_RO_RO11219082)分组,然后获取每个分组中最后一段数字最小和最大的完整记录。之前的尝试只得到了分组键和数字本身,没有保留完整的字段,这是因为你用了MIN()/MAX()聚合单个字段,而不是获取对应条件的完整行。
正确实现方案
KSQL(你使用的语法是KSQL)提供了MIN_BY()和MAX_BY()函数,专门用于返回满足最小/最大排序键条件的完整行或指定字段。我们可以利用这两个函数来实现你的需求:
1. 获取每个分组的最小值对应完整记录
CREATE TABLE HISTORY_MIN WITH (KAFKA_TOPIC='history_min', PARTITIONS=10, REPLICAS=1, FORMAT='JSON') AS SELECT -- 构造分组键:将payload_1的前两部分(history + 唯一标识段)拼接 CONCAT_WS('::', SPLIT(my_key_col->`payload`, '::')[0], SPLIT(my_key_col->`payload`, '::')[1]) AS group_key, -- 获取数字最小的payload_1 MIN_BY(my_key_col->`payload`, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS payload_1, -- 同步获取对应记录的schema字段 MIN_BY(`schema`, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS schema, -- 同步获取对应记录的payload字段 MIN_BY(payload, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS payload FROM HISTORY_CONTENT GROUP BY CONCAT_WS('::', SPLIT(my_key_col->`payload`, '::')[0], SPLIT(my_key_col->`payload`, '::')[1]);
2. 获取每个分组的最大值对应完整记录
CREATE TABLE HISTORY_MAX WITH (KAFKA_TOPIC='history_max', PARTITIONS=10, REPLICAS=1, FORMAT='JSON') AS SELECT CONCAT_WS('::', SPLIT(my_key_col->`payload`, '::')[0], SPLIT(my_key_col->`payload`, '::')[1]) AS group_key, MAX_BY(my_key_col->`payload`, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS payload_1, MAX_BY(`schema`, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS schema, MAX_BY(payload, CAST(SPLIT(my_key_col->`payload`, '::')[2] AS BIGINT)) AS payload FROM HISTORY_CONTENT GROUP BY CONCAT_WS('::', SPLIT(my_key_col->`payload`, '::')[0], SPLIT(my_key_col->`payload`, '::')[1]);
关键说明
- 分组键构造:用
CONCAT_WS('::', ...)把payload_1拆分后的前两部分拼接,确保同一个业务标识下的记录被分到同一组。 - 类型转换:必须把
payload_1最后一段的字符串转成BIGINT,否则会按字符串字典序排序(比如"10"会被认为比"2"小),导致结果错误。 - MIN_BY/MAX_BY的作用:这两个函数会根据第二个参数(排序键)的最小/最大值,返回第一个参数对应的字段值,完美匹配你需要获取完整记录的需求。
验证查询
创建完成后,你可以用以下语句验证结果:
-- 查看最小值记录 SELECT payload_1, schema, payload FROM HISTORY_MIN; -- 查看最大值记录 SELECT payload_1, schema, payload FROM HISTORY_MAX;
内容的提问来源于stack exchange,提问作者sbrienne
相关产品推荐
相关产品推荐

