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

如何通过KSQL实现无Schema的Kafka消息过滤复制及JSON格式转换

KSQL过滤JSON消息后保留原始格式的解决方案

方法一:基于JSON路径提取过滤字段(推荐)

直接利用KSQL的JSON路径支持,无需定义完整Schema,同时保留原始消息结构:

  1. 创建源流,仅指定需要过滤的字段路径
CREATE STREAM Test1 (
  filter_attr VARCHAR PATH '$.your_target_attribute' -- 替换为你的过滤属性路径
) WITH (
  KAFKA_TOPIC='your_source_topic',
  VALUE_FORMAT='JSON',
  KEY_FORMAT='KAFKA'
);
  1. 过滤并生成目标流,直接输出原始消息
CREATE STREAM Test2 WITH (KAFKA_TOPIC='your_target_topic') AS
SELECT VALUE -- 直接选择原始JSON消息
FROM Test1
WHERE filter_attr = 'A';

方法二:兼容已有的VARCHAR类型流

如果已经将整个消息定义为VARCHAR类型的data字段,通过类型转换恢复JSON结构:

CREATE STREAM Test2 WITH (
  KAFKA_TOPIC='your_target_topic',
  VALUE_FORMAT='JSON' -- 指定输出格式为JSON
) AS
SELECT CAST(data AS JSON) AS VALUE
FROM Test1
WHERE EXTRACTJSONFIELD(data, '$.your_target_attribute') = 'A';

关键注意事项

  • 目标流必须指定VALUE_FORMAT='JSON'(方法一中因继承源流格式可省略,但显式指定更稳妥)
  • 避免选择单个字段输出,直接选择原始VALUE或转换后的JSON类型字段,才能保留完整原始消息结构

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 17:18:26