使用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
相关产品推荐
相关产品推荐

