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

如何解析ksql中Tumbling滚动窗口输出的JSON格式消息键

KSQL Tumbling窗口Key反序列化解决方案

方案1:调整KSQL输出规则(推荐,无需额外解析逻辑)

操作步骤:

  • 在执行建表语句前,先设置KSQL配置参数:
SET ksql.windowed.key.format = JSON;
SET ksql.output.windowed.key.separator = '|';
  • 若要长期生效可直接在ksql-server.properties中添加如下配置,避免每次执行SQL都手动设置:
ksql.windowed.key.format=JSON
ksql.output.windowed.key.separator=|
# 可选配置:输出窗口结束时间戳,默认关闭时Tumbling窗口结束位显示为-
ksql.output.windowed.key.include.end=true

调整后输出的窗口键会被拆分为业务分组字段JSON|窗口开始时间戳|窗口结束时间戳的格式,仅需要按分隔符分割字符串,分别取对应部分即可直接使用:

调整后key样例:{"ACCOUNTID":"account1","DEVICEID":"device27","ENTITY":"Account"}|1636090860000|-

方案2:消费者端直接解析原生窗口键

如果不方便调整KSQL配置,可以直接按固定规则解析原生key字符串:
原生窗口键格式固定为:[{业务分组JSON}]@{窗口起始时间戳}/{窗口结束时间戳}

  • 解析逻辑步骤:
  1. 去除key首尾的方括号[、]
  2. 按@字符分割为两部分,前半部分就是业务分组字段的JSON字符串,可直接反序列化
  3. 后半部分按/分割,第一部分是窗口起始时间戳(毫秒级Long类型),第二部分是窗口结束时间戳(Tumbling窗口未配置输出结束时间时此处为-,可以按窗口大小=起始时间+60*1000计算得到结束时间)

Java语言解析示例代码:

import com.alibaba.fastjson.JSONObject;

public class KsqlWindowKeyParser {
    public static void parse(String windowKey) {
        // 去除首尾[]
        String content = windowKey.substring(1, windowKey.length() - 1);
        // 分割业务JSON和窗口时间部分
        String[] parts = content.split("@", 2);
        String groupJson = parts[0];
        String windowPart = parts[1];
        // 解析分组字段
        JSONObject groupObj = JSONObject.parseObject(groupJson);
        String accountId = groupObj.getString("ACCOUNTID");
        String deviceId = groupObj.getString("DEVICEID");
        String entity = groupObj.getString("ENTITY");
        // 解析窗口时间
        String[] windowTimes = windowPart.split("/", 2);
        Long windowStart = Long.parseLong(windowTimes[0]);
        Long windowEnd = windowStart + 60 * 1000;
    }
}

Python语言解析示例代码:

import json

def parse_ksql_window_key(window_key: str):
    # 去除首尾[]
    content = window_key.strip("[]")
    # 分割业务JSON和窗口部分
    group_json, window_part = content.split("@", 1)
    # 解析分组字段
    group_data = json.loads(group_json)
    account_id = group_data["ACCOUNTID"]
    device_id = group_data["DEVICEID"]
    entity = group_data["ENTITY"]
    # 解析窗口时间
    window_start_str, window_end_str = window_part.split("/", 1)
    window_start = int(window_start_str)
    window_end = window_start + 60 * 1000
    return group_data, window_start, window_end

内容的提问来源于stack exchange,提问作者Ajay Sinha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 01:27:08