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

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数组的类类型标识)类型的值,无法完成强转触发序列化失败。

根因分析

异常由配置冲突和代码逻辑问题共同触发:

  1. 序列化逻辑冲突:使用FlinkKafkaProducer010时传入了SimpleStringSchema,Flink会在算子内部先把String类型的数据序列化为byte数组,再传给底层Kafka客户端。但代码中手动在Properties里配置了Kafka原生的StringSerializer,底层客户端收到Flink传入的byte数组后,会尝试将其强转为String再做序列化,直接触发ClassCastException。
  2. JSON输出逻辑错误:Map算子最后调用obj.asText()返回结果,对于Object类型的JsonNode,该方法不会返回完整的JSON结构字符串,仅会返回节点的文本值,不符合输出完整扩充后JSON的业务预期。
  3. 冗余配置: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 19:56:49