Mongo Sink Connector启动失败:如何用ksql生成Map/Struct格式记录键?
解决Mongo Sink Connector启动失败及KSql生成Struct/Map类型键的方法
问题原因
你使用的FullKeyStrategy文档ID策略要求Kafka记录的键必须是Map或Struct类型,但通过PARTITION BY单个字段生成的键是字符串等基础类型,不符合要求,导致启动失败。
方案一:通过KSql生成Struct/Map类型的记录键
生成Struct类型键
如果需要将多个字段组合成Struct作为键,或者把单个字段包装成Struct,可以用STRUCT()函数构造:
-- 多字段组合成Struct键 CREATE STREAM target_stream_struct_key AS SELECT STRUCT(user_id := user_id, order_no := order_no) AS record_key, * EXCLUDE(user_id, order_no) -- 排除已放入键的字段(可选) FROM source_stream PARTITION BY record_key; -- 单个字段包装成Struct键 CREATE STREAM target_stream_single_struct AS SELECT STRUCT(id := user_id) AS record_key, * EXCLUDE(user_id) FROM source_stream PARTITION BY record_key;
生成Map类型键
使用MAP()函数将字段键值对组合成Map作为记录键:
CREATE STREAM target_stream_map_key AS SELECT MAP(ARRAY['user_id', 'order_no'], ARRAY[user_id, order_no]) AS record_key, * EXCLUDE(user_id, order_no) FROM source_stream PARTITION BY record_key;
方案二:调整Mongo Sink Connector配置(无需修改KSql键类型)
如果不想修改KSql的键生成逻辑,可以更换文档ID策略,改用数据中的字段作为Mongo的_id:
- 修改Connector配置:
Document ID strategy: ProvidedInValueStrategy Document ID strategy field: mongo_document_id - 在KSql中生成用于Mongo
_id的字段:
这样Connector会直接从消息value中读取CREATE STREAM source_stream_with_id AS SELECT -- 可以用字段拼接、UUID()等方式生成唯一ID user_id || '_' || order_no AS mongo_document_id, * FROM source_stream;mongo_document_id作为Mongo文档的ID,不再依赖记录键的类型。
内容的提问来源于stack exchange,提问作者sai jyothsna pentyala
相关产品推荐
相关产品推荐

