Apache Flink中无法将数据写入KafkaSink问题求助
Flink KafkaSink写入时出现Kafka节点断开问题
我用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
相关产品推荐
相关产品推荐

