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

使用STRUCT类型创建KSQL Stream失败,请求技术协助

问题解决:KSQL使用STRUCT类型读取MongoDB同步的Kafka主题

问题根源

从Kafka主题的print输出可以看到,消息value中的payload字段是转义后的JSON字符串,而非原生的JSON对象。你之前直接将payload定义为STRUCT类型时,KSQL会尝试把字符串当作JSON对象解析,解析失败后返回空值。

解决方案

需要先读取原始的JSON结构,再将payload字符串解析为STRUCT类型,分两步实现:

1. 创建原始数据流

先定义匹配Kafka消息实际结构的流,将payload按字符串类型读取:

CREATE STREAM RAW_SUPERSTORE_PEOPLE (
  schema STRUCT<type VARCHAR, optional BOOLEAN>,
  payload VARCHAR
) WITH (
  KAFKA_TOPIC='Mongo.Sample_SuperStore.People',
  VALUE_FORMAT='JSON'
);

2. 解析字符串为STRUCT并创建派生流

使用FROM_JSON函数将payload字符串解析为目标STRUCT结构,生成新的流:

CREATE STREAM STREAM_SUPERSTORE_PEOPLE AS
SELECT
  FROM_JSON(
    payload,
    'STRUCT<_id:STRUCT<`$oid`:VARCHAR>, CustomerID:VARCHAR, CustomerName:VARCHAR, Segment:VARCHAR, Country:VARCHAR, City:VARCHAR, State:VARCHAR, PostalCode:INT, Region:VARCHAR>'
  ) AS customer
FROM RAW_SUPERSTORE_PEOPLE
EMIT CHANGES;

3. 查询验证

现在可以直接查询STRUCT中的字段:

SELECT 
  customer._id.`$oid` AS oid,
  customer.CustomerID,
  customer.CustomerName,
  customer.Segment
FROM STREAM_SUPERSTORE_PEOPLE
EMIT CHANGES LIMIT 1;

替代方案:直接在查询中解析

如果不需要持久化派生流,也可以直接在查询时解析字符串:

SELECT
  FROM_JSON(
    payload,
    'STRUCT<_id:STRUCT<`$oid`:VARCHAR>, CustomerID:VARCHAR, CustomerName:VARCHAR, Segment:VARCHAR, Country:VARCHAR, City:VARCHAR, State:VARCHAR, PostalCode:INT, Region:VARCHAR>'
  ).*
FROM RAW_SUPERSTORE_PEOPLE
EMIT CHANGES LIMIT 1;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 12:18:25