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

Scala开发Flink任务消费Base64编码Kafka消息如何解码为有效JSON

问题原因

你遇到的报错和乱码本质是处理流程和数据格式不匹配:

  • 第一次直接用JsonNodeDeserializationSchema失败,是因为Kafka中存储的不是原生JSON字节,是Base64编码后的压缩二进制,不符合JSON格式校验
  • 第二次用SimpleStringSchema拿到乱码,是因为该Schema直接把压缩后的二进制字节按照UTF-8编码转成字符串,压缩字节本身不符合UTF-8编码规则,就会出现乱码字符
解决方案

核心处理流程为:获取Kafka原始字节 → Base64解码 → 按生产端约定的压缩算法解压 → 转成JSON格式,你可以直接自定义反序列化Schema实现全流程处理,无需额外单独加算子转换。

步骤1:实现自定义反序列化Schema

import java.io.ByteArrayInputStream
import java.util.Base64
import java.util.zip.GZIPInputStream
import com.fasterxml.jackson.databind.node.ObjectNode
import com.fasterxml.jackson.databind.ObjectMapper
import org.apache.flink.api.common.serialization.DeserializationSchema
import org.apache.flink.api.common.typeinfo.TypeInformation
import org.apache.commons.io.IOUtils
import org.apache.flink.api.java.typeutils.TypeExtractor

class Base64CompressedJsonDeserializationSchema(compressType: String = "gzip") extends DeserializationSchema[ObjectNode] {
  // 懒加载避免序列化问题
  @transient private lazy val objectMapper = new ObjectMapper()
  @transient private lazy val base64Decoder = Base64.getDecoder

  override def deserialize(message: Array[Byte]): ObjectNode = {
    // 1. Base64解码原始字节
    val decodedBytes = base64Decoder.decode(message)
    // 2. 按指定压缩格式解压
    val decompressedStr = compressType match {
      case "gzip" =>
        val gzipIs = new GZIPInputStream(new ByteArrayInputStream(decodedBytes))
        IOUtils.toString(gzipIs, "UTF-8")
      case "snappy" =>
        // 使用snappy需提前引入snappy-java依赖
        org.xerial.snappy.Snappy.uncompressString(decodedBytes, "UTF-8")
      case "deflate" =>
        // deflate格式解压逻辑可按需补充
        val inflater = new java.util.zip.Inflater()
        inflater.setInput(decodedBytes)
        val resultBytes = new Array[Byte](1024 * 1024)
        val length = inflater.inflate(resultBytes)
        inflater.end()
        new String(resultBytes, 0, length, "UTF-8")
      case _ =>
        // 无压缩场景直接转字符串
        new String(decodedBytes, "UTF-8")
    }
    // 3. 转换为ObjectNode返回
    objectMapper.readTree(decompressedStr).asInstanceOf[ObjectNode]
  }

  override def isEndOfStream(nextElement: ObjectNode): Boolean = false

  override def getProducedType: TypeInformation[ObjectNode] = TypeExtractor.getForClass(classOf[ObjectNode])
}

步骤2:替换Kafka消费者的反序列化Schema

val x_stream: DataStream[ObjectNode] = env
  .addSource(
    new FlinkKafkaConsumer010[ObjectNode](
      parameters.get("kafka.topic.source"),
      // 压缩类型需和Kafka生产端配置保持一致,示例为gzip
      new Base64CompressedJsonDeserializationSchema("gzip"),
      kfk_props
    ).setStartFromEarliest()
  ).name("Topic Test")
  .rebalance
注意事项
  • 如果使用snappy压缩,需要在项目依赖中引入snappy-java包,gzip、deflate为JDK内置能力无需额外依赖
  • 可以在deserialize方法中增加try-catch逻辑,捕获解码、解压、JSON转换异常,过滤脏数据避免任务崩溃
  • 压缩算法必须和生产端写入Kafka时的压缩算法完全匹配,否则会解压失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 13:48:00