自定义Kafka Connect值转换器输出仍为Base64问题求助
分析与解决你的Kafka Connect值转换器问题
看起来你遇到的核心问题是JsonConverter处理字节数组时默认会序列化为Base64格式,这直接导致了你消费输出主题时看到的还是Base64内容。我帮你拆解下代码里的问题,并给出修改方案:
问题根源
你的代码中,当decode=true时,把ByteArrayConverter解码后的字节数组直接传给了JsonConverter.fromConnectData。但JsonConverter对字节数组类型的value,默认行为是将其转为Base64编码的字符串,而不是把字节数组当成JSON字符串去解析。这就是为什么你打印解码后的字符串是正常的,但最终输出还是Base64格式。
另外还有两个小问题:
- 代码中重复调用了
decoder.fromConnectData,做了冗余操作 - 没有在
configure方法中初始化decoder和delegate转换器,可能导致配置不生效
修改后的代码示例
private final Converter delegate = new JsonConverter(); private final Converter decoder = new ByteArrayConverter(); private boolean decode = false; private final ObjectMapper objectMapper = new ObjectMapper(); // 用于JSON合法性校验与解析 @Override public void configure(Map<String, ?> configs, boolean isKey) { // 必须初始化两个转换器,传入配置参数 decoder.configure(configs, isKey); delegate.configure(configs, isKey); // 从配置中读取decode开关,也可以直接设为true decode = Boolean.parseBoolean(configs.getOrDefault("decode", "true").toString()); } @Override public byte[] fromConnectData(String topic, Schema schema, Object value) { try { // 只解码一次,避免重复操作 byte[] decodedBytes = decoder.fromConnectData(topic, schema, value); String decodedString = new String(decodedBytes, StandardCharsets.UTF_8); LOG.info("decoded string: {}", decodedString); if (decode) { // 先把解码后的字符串解析为JSON对象,再传给JsonConverter Object jsonObject = objectMapper.readValue(decodedString, Object.class); return delegate.fromConnectData(topic, schema, jsonObject); } else { // 不需要解码时直接委托处理 return delegate.fromConnectData(topic, schema, value); } } catch (Exception e) { LOG.error("Failed to process record", e); // 包装无效JSON为合法的JSON对象 String invalidContent = new String(decoder.fromConnectData(topic, schema, value), StandardCharsets.UTF_8); String wrappedJson = wrapInvalidJson(invalidContent); try { // 将包装后的JSON转为对象再传给转换器 Object wrappedObject = objectMapper.readValue(wrappedJson, Object.class); return delegate.fromConnectData(topic, schema, wrappedObject); } catch (JsonProcessingException ex) { // 极端情况:包装后的JSON也异常,直接返回字节数组 LOG.error("Failed to process wrapped JSON", ex); return wrappedJson.getBytes(StandardCharsets.UTF_8); } } } // 示例:将无效JSON包装为合法的JSON对象 private String wrapInvalidJson(String invalidJson) throws JsonProcessingException { Map<String, String> wrapper = new HashMap<>(); wrapper.put("invalidJson", invalidJson); return objectMapper.writeValueAsString(wrapper); }
关键修改点说明
- 初始化转换器:在
configure方法中调用两个转换器的configure方法,确保它们加载正确的配置(比如ByteArrayConverter需要启用Base64解码)。 - 避免重复解码:只调用一次
decoder.fromConnectData,减少冗余操作。 - 先解析JSON再转换:用
ObjectMapper把解码后的字符串解析为Java对象(Map/List等),再传给JsonConverter,这样转换器会生成正常的JSON结构,而不是Base64。 - 明确字符编码:指定
UTF_8编码,避免因系统默认编码不同导致的乱码问题。
额外注意事项
- 确保
ByteArrayConverter的配置正确:因为你从Kinesis读取的是Base64数据,需要给ByteArrayConverter配置bytearray.converter.decoder=base64。 - 测试边界场景:比如空字符串、完全不符合JSON格式的内容,验证
wrapInvalidJson是否能正确生成合法JSON。
内容的提问来源于stack exchange,提问作者Oded Rosenberg
相关产品推荐
相关产品推荐

