diff --git a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java index a3a8aec5..c4605067 100644 --- a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java @@ -1,5 +1,6 @@ package com.alibaba.otter.canal.kafka; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Properties; @@ -138,7 +139,7 @@ public class CanalKafkaProducer implements CanalMQProducer { private void send(MQProperties.CanalDestination canalDestination, String topicName, Message message) throws Exception { if (!kafkaProperties.getFlatMessage()) { - ProducerRecord record = null; + List> records = new ArrayList<>(); if (canalDestination.getPartitionHash() != null && !canalDestination.getPartitionHash().isEmpty()) { Message[] messages = MQMessageUtils.messagePartition(message, canalDestination.getPartitionsNum(), @@ -147,25 +148,23 @@ public class CanalKafkaProducer implements CanalMQProducer { for (int i = 0; i < length; i++) { Message messagePartition = messages[i]; if (messagePartition != null) { - record = new ProducerRecord<>(topicName, i, null, messagePartition); + records.add(new ProducerRecord<>(topicName, i, null, messagePartition)); } } } else { final int partition = canalDestination.getPartition() != null ? canalDestination.getPartition() : 0; - record = new ProducerRecord<>(topicName, partition, null, message); + records.add(new ProducerRecord<>(topicName, partition, null, message)); } - if (record != null) { - if (kafkaProperties.getTransaction()) { - producer.send(record).get(); - } else { - producer.send(record).get(); - } - - if (logger.isDebugEnabled()) { - logger.debug("Send message to kafka topic: [{}], packet: {}", topicName, message.toString()); - } - } + if (!records.isEmpty()) { + for (ProducerRecord record : records) { + producer.send(record).get(); + } + + if (logger.isDebugEnabled()) { + logger.debug("Send message to kafka topic: [{}], packet: {}", topicName, message.toString()); + } + } } else { // 发送扁平数据json List flatMessages = MQMessageUtils.messageConverter(message);