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
相关产品推荐
相关产品推荐

