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

Apache Flink中无法将数据写入KafkaSink问题求助

我用Java编写的Apache Flink程序,从Kafka读取{name: "abc", age: 20}格式的数据,尝试将数据写回Kafka,但调用data.sinkTo(sink)后出现Kafka节点断开的问题,相关信息如下:

KafkaSink初始化代码

KafkaSink<AllIncidentsDataPOJO> sink = KafkaSink.<AllIncidentsDataPOJO>builder()
                        .setBootstrapServers(this.bootstrapServer)
                        .setKafkaProducerConfig(kafkaProps)
                        .setRecordSerializer(KafkaRecordSerializationSchema.builder()
                                        .setTopic("flink-all-incidents")
                                        .setValueSerializationSchema(new IncidentsSerializationSchema()).build())
                        .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
                        .build();

IncidentsSerializationSchema实现类

public class IncidentsSerializationSchema implements SerializationSchema<AllIncidentsDataPOJO> {
    static ObjectMapper objectMapper = new ObjectMapper();
    @Override
    public byte[] serialize(AllIncidentsDataPOJO element) {
        if (objectMapper == null) {
            objectMapper.setVisibility(PropertyAccessor.FIELD, JsonAutoDetect.Visibility.ANY);
            objectMapper = new ObjectMapper();
        }
        try {
//            System.out.println("Returned value: " + objectMapper.writeValueAsString(element).getBytes());
            return objectMapper.writeValueAsString(element).getBytes();
        } catch (
                org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException e) {
            System.out.println("Exception: " + e);
        }
        return new byte[0];
    }
}

序列化打印信息

序列化时打印的字节数组标识:

Returned value: [B@6f29a884
Returned value: [B@743145c5
Returned value: [B@6aa6b431
Returned value: [B@208c0151
Returned value: [B@1a3a0c6d
Returned value: [B@7972da35
Returned value: [B@1097059f
Returned value: [B@6013a87b
Returned value: [B@27b252dc
Returned value: [B@288478b7
Returned value: [B@54185041
Returned value: [B@970646b

报错日志

[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Node 8 disconnected.
[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Cancelled in-flight API_VERSIONS request with correlation id 35 due to node 8 being disconnected (elapsed time since creation: 292ms, elapsed time since send: 292ms, request timeout: 30000ms)
[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Node 8 disconnected.
[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Cancelled in-flight API_VERSIONS request with correlation id 37 due to node 8 being disconnected (elapsed time since creation: 639ms, elapsed time since send: 639ms, request timeout: 30000ms)
[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Node 8 disconnected.
[kafka-producer-network-thread | producer-kafka-sink-0-1] INFO org.apache.kafka.clients.NetworkClient - [Producer clientId=producer-kafka-sink-0-1, transactionalId=kafka-sink-0-1] Cancelled in-flight API_VERSIONS request with correlation id 38 due to node 8 being disconnected (elapsed time since creation: 875ms, elapsed time since send: 875ms, request timeout: 30000ms)

可能的解决方案

1. 修复序列化类的线程安全及逻辑问题

当前IncidentsSerializationSchema的静态ObjectMapper初始化逻辑存在错误:静态变量已初始化,if (objectMapper == null)分支永远不会执行,导致配置的字段可见性不生效;同时序列化失败时返回空字节数组会导致静默错误。修改后的代码:

public class IncidentsSerializationSchema implements SerializationSchema<AllIncidentsDataPOJO> {
    private static final ObjectMapper objectMapper;

    static {
        objectMapper = new ObjectMapper();
        objectMapper.setVisibility(PropertyAccessor.FIELD, JsonAutoDetect.Visibility.ANY);
    }

    @Override
    public byte[] serialize(AllIncidentsDataPOJO element) {
        try {
            return objectMapper.writeValueAsBytes(element); // 直接使用更高效的writeValueAsBytes方法
        } catch (org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException e) {
            throw new RuntimeException("序列化AllIncidentsDataPOJO失败", e); // 抛出异常暴露问题
        }
    }
}

2. 检查Exactly-Once投递语义的配置正确性

开启EXACTLY_ONCE需要满足:

  • Kafka集群版本≥0.11.0.0(支持事务)
  • Kafka Broker配置开启事务:transaction.state.log.replication.factor≥3,transaction.state.log.min.isr≥2
  • 确保Flink集群与Kafka集群的时间同步

3. 排查网络与Kafka节点状态

  • 确认Flink集群能访问Kafka节点8的端口(默认9092/9093)
  • 检查Kafka节点8是否正常运行,是否存在磁盘满、网络波动等问题
  • 临时将投递语义改为AT_LEAST_ONCE,验证是否能正常写入,判断是否为事务配置问题

4. 调整Kafka Producer重试与重连配置

在kafkaProps中添加以下配置,提升网络容错能力:

kafkaProps.put(ProducerConfig.RETRIES_CONFIG, 3);
kafkaProps.put(ProducerConfig.RECONNECT_BACKOFF_MS_CONFIG, 1000);
kafkaProps.put(ProducerConfig.RECONNECT_BACKOFF_MAX_MS_CONFIG, 5000);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 22:01:03