如何为Kafka Connect的SMT传递列表嵌套列表配置?
解决Kafka Connect SMT嵌套列表配置的反序列化错误
问题根源
Kafka Connect的ConfigDef.Type.LIST仅支持一维逗号分隔的字符串列表(比如"field1,field2"),无法直接处理嵌套数组的JSON语法。本地测试时你可能直接通过Java Map传入了List对象,所以解析正常;但通过REST API部署时,配置被序列化为JSON,REST层尝试将嵌套数组转换为String类型失败,触发反序列化错误。
解决方案
方案1:将嵌套列表配置改为JSON字符串(推荐)
1. 调整配置格式
把原配置中的嵌套数组改为转义后的JSON字符串:
"transforms": "EqualityCheck", "transforms.EqualityCheck.type":"org.apache.kafka.connect.transforms.EqualityCheckOnFields$Value", "transforms.EqualityCheck.fields.notEquality": "field1, field2", "transforms.EqualityCheck.fieldValues.notEquality": "[[9, 0, null, 4.8234, 9.9999], [\"pol\", \"POC\", \"\", null]]"
2. 修改ConfigDef定义
将fieldValues.notEquality的类型从LIST改为STRING:
public static final ConfigDef CONFIG_DEF = new ConfigDef() .define("fields.notEquality", ConfigDef.Type.LIST, Collections.emptyList(), ConfigDef.Importance.MEDIUM, "List of fields for not-equality check") .define("fieldValues.notEquality", ConfigDef.Type.STRING, "", ConfigDef.Importance.MEDIUM, "JSON string representing list of lists containing invalid values");
3. 简化配置解析逻辑
直接将JSON字符串解析为嵌套列表,无需处理混合类型:
@Override public void configure(Map<String, ?> props) { SimpleConfig config = new SimpleConfig(CONFIG_DEF, props); notEqualityFields = config.getList("fields.notEquality"); try { ObjectMapper mapper = new ObjectMapper(); // 使用TypeReference明确嵌套列表类型 notEqualityValues = mapper.readValue( config.getString("fieldValues.notEquality"), new TypeReference<List<List<?>>>() {} ); } catch (IOException e) { throw new ConfigException("Invalid JSON format for fieldValues.notEquality", e); } }
方案2:保留LIST类型,用逗号分隔JSON字符串
如果坚持使用ConfigDef.Type.LIST,可以将每个嵌套数组作为独立的JSON字符串,用逗号分隔:
1. 调整配置格式
"transforms.EqualityCheck.fieldValues.notEquality": "[9, 0, null, 4.8234, 9.9999], [\"pol\", \"POC\", \"\", null]"
2. 保留原解析逻辑
你的parseValues方法可以直接处理这种格式,因为每个元素都是JSON字符串,会被解析为List对象。
验证
部署修改后的连接器,REST层会将配置值作为字符串传递,你的代码可以正确解析为嵌套列表,避免反序列化错误。
内容的提问来源于stack exchange,提问作者Priyanshu Sharma
相关产品推荐
相关产品推荐

