Flink向Kafka发布消息出现SerializationException序列化错误
Flink写入Kafka序列化异常排查方案
问题复现代码
以下是触发异常的业务逻辑代码:
final Properties props = new Properties(); props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put("security.protocol", "SASL_SSL"); props.put("sasl.mechanism", "PLAIN"); props.put("client.dns.lookup", "use_all_dns_ips"); props.put("group.id", "flink-consumer-1"); DataStream<String> enrichedJson = flattenedJson.map(new MapFunction<JsonNode, String>() { public String map(JsonNode value) throws Exception { Integer customerId = value.get("customerId").intValue(); ObjectMapper mapper = new ObjectMapper(); String json = mapper.writeValueAsString(map.get(customerId).get("name")); JsonNode jsonNode = mapper.readTree(json); JsonNode obj = ((ObjectNode) value).set("name", jsonNode ); json = mapper.writeValueAsString(map.get(customerId).get("mobileNumber")); jsonNode = mapper.readTree(json); obj = ((ObjectNode) value).set("mobileNumber", jsonNode ); return obj.asText(); } }); enrichedJson.print(); FlinkKafkaProducer010 <String> producer = new FlinkKafkaProducer010<String>("ENRICHED_CUSTOMER", new SimpleStringSchema(), props); enrichedJson.addSink(producer);
异常信息
作业运行时抛出序列化异常,核心堆栈如下:
Exception in thread "main" org.apache.flink.runtime.client.JobExecutionException: org.apache.kafka.common.errors.SerializationException: Can't convert value of class [B to class org.apache.kafka.common.serialization.StringSerializer specified in value.serializer at org.apache.flink.runtime.minicluster.MiniCluster.executeJobBlocking(MiniCluster.java:625) at org.apache.flink.streaming.api.environment.LocalStreamEnvironment.execute(LocalStreamEnvironment.java:121) at flink.KafkaFlink.main(KafkaFlink.java:123) Caused by: org.apache.kafka.common.errors.SerializationException: Can't convert value of class [B to class org.apache.kafka.common.serialization.StringSerializer specified in value.serializer at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:955) at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:912) at org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer010.invoke(FlinkKafkaProducer010.java:382) ... 省略部分堆栈 Caused by: java.lang.ClassCastException: class [B cannot be cast to class java.lang.String ([B and java.lang.String are in module java.base of loader 'bootstrap') at org.apache.kafka.common.serialization.StringSerializer.serialize(StringSerializer.java:29) at org.apache.kafka.common.serialization.Serializer.serialize(Serializer.java:62) at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:952) ... 43 more
核心错误点:Kafka生产者配置
StringSerializer作为值序列化器,但实际接收到字节数组([B是Java中byte数组的类类型标识)类型的值,无法完成强转触发序列化失败。
根因分析
异常由配置冲突和代码逻辑问题共同触发:
- 序列化逻辑冲突:使用
FlinkKafkaProducer010时传入了SimpleStringSchema,Flink会在算子内部先把String类型的数据序列化为byte数组,再传给底层Kafka客户端。但代码中手动在Properties里配置了Kafka原生的StringSerializer,底层客户端收到Flink传入的byte数组后,会尝试将其强转为String再做序列化,直接触发ClassCastException。 - JSON输出逻辑错误:Map算子最后调用
obj.asText()返回结果,对于Object类型的JsonNode,该方法不会返回完整的JSON结构字符串,仅会返回节点的文本值,不符合输出完整扩充后JSON的业务预期。 - 冗余配置:
group.id是Kafka消费者专属配置,在生产者配置中设置不生效,属于无效配置。
修复方案
按照以下步骤修改代码即可解决问题:
- 删除Properties中冗余的Kafka原生序列化器配置,让Flink的SerializationSchema完全接管序列化逻辑:
// 删除以下两行冗余配置,不要给FlinkKafkaProducer手动指定Kafka原生key/value序列化器 // props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); // props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); // 可选删除无效的消费者配置 // props.put("group.id", "flink-consumer-1");
- 修正Map算子的JSON处理逻辑,将ObjectMapper实例提到算子外部避免每条数据重复创建(减少性能开销),替换错误的
asText()返回逻辑为正确的JSON序列化方式,同时简化冗余的序列化/反序列化步骤:
// 初始化ObjectMapper放在算子外部,单实例复用 ObjectMapper mapper = new ObjectMapper(); DataStream<String> enrichedJson = flattenedJson.map(new MapFunction<JsonNode, String>() { @Override public String map(JsonNode value) throws Exception { Integer customerId = value.get("customerId").intValue(); ObjectNode node = (ObjectNode) value; // 直接从维度关联的map中取值写入节点,不需要重复做writeValueAsString/readTree的冗余操作 node.put("name", map.get(customerId).get("name").asText()); node.put("mobileNumber", map.get(customerId).get("mobileNumber").asText()); // 正确输出完整JSON结构的字符串 return mapper.writeValueAsString(node); } });
修改完成后重新提交作业,序列化异常会消失,Kafka中可以正常写入扩充后的完整JSON数据。
内容的提问来源于stack exchange,提问作者user1112259
相关产品推荐
相关产品推荐

