You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何获取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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 15:54:04