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

自定义Kafka Connect值转换器输出仍为Base64问题求助

分析与解决你的Kafka Connect值转换器问题

看起来你遇到的核心问题是JsonConverter处理字节数组时默认会序列化为Base64格式,这直接导致了你消费输出主题时看到的还是Base64内容。我帮你拆解下代码里的问题,并给出修改方案:

问题根源

你的代码中,当decode=true时,把ByteArrayConverter解码后的字节数组直接传给了JsonConverter.fromConnectData。但JsonConverter对字节数组类型的value,默认行为是将其转为Base64编码的字符串,而不是把字节数组当成JSON字符串去解析。这就是为什么你打印解码后的字符串是正常的,但最终输出还是Base64格式。

另外还有两个小问题:

  1. 代码中重复调用了decoder.fromConnectData,做了冗余操作
  2. 没有在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);
}

关键修改点说明

  1. 初始化转换器:在configure方法中调用两个转换器的configure方法,确保它们加载正确的配置(比如ByteArrayConverter需要启用Base64解码)。
  2. 避免重复解码:只调用一次decoder.fromConnectData,减少冗余操作。
  3. 先解析JSON再转换:用ObjectMapper把解码后的字符串解析为Java对象(Map/List等),再传给JsonConverter,这样转换器会生成正常的JSON结构,而不是Base64。
  4. 明确字符编码:指定UTF_8编码,避免因系统默认编码不同导致的乱码问题。

额外注意事项

  • 确保ByteArrayConverter的配置正确:因为你从Kinesis读取的是Base64数据,需要给ByteArrayConverter配置bytearray.converter.decoder=base64。
  • 测试边界场景:比如空字符串、完全不符合JSON格式的内容,验证wrapInvalidJson是否能正确生成合法JSON。

内容的提问来源于stack exchange,提问作者Oded Rosenberg

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 07:02:34