浏览代码

kafka单个partition的发送未改为同步发送的bug fix

mcy 6 年之前
父节点
当前提交
b71e299161
共有 1 个文件被更改,包括 1 次插入1 次删除
  1. 1 1
      server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java

+ 1 - 1
server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java

@@ -105,7 +105,7 @@ public class CanalKafkaProducer implements CanalMQProducer {
                                 canalDestination.getPartition(),
                                 null,
                                 JSON.toJSONString(flatMessage));
-                            producer2.send(record);
+                            producer2.send(record).get();
                         } catch (Exception e) {
                             logger.error(e.getMessage(), e);
                             // producer.abortTransaction();