Bladeren bron

修改接口引用

guixincui 6 jaren geleden
bovenliggende
commit
90dab6831b
1 gewijzigde bestanden met toevoegingen van 1 en 2 verwijderingen
  1. 1 2
      server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java

+ 1 - 2
server/src/main/java/com/alibaba/otter/canal/server/CanalMQStarter.java

@@ -13,7 +13,6 @@ import org.slf4j.LoggerFactory;
 import com.alibaba.otter.canal.common.MQProperties;
 import com.alibaba.otter.canal.common.MQProperties;
 import com.alibaba.otter.canal.instance.core.CanalInstance;
 import com.alibaba.otter.canal.instance.core.CanalInstance;
 import com.alibaba.otter.canal.instance.core.CanalMQConfig;
 import com.alibaba.otter.canal.instance.core.CanalMQConfig;
-import com.alibaba.otter.canal.kafka.CanalKafkaProducer;
 import com.alibaba.otter.canal.protocol.ClientIdentity;
 import com.alibaba.otter.canal.protocol.ClientIdentity;
 import com.alibaba.otter.canal.protocol.Message;
 import com.alibaba.otter.canal.protocol.Message;
 import com.alibaba.otter.canal.server.embedded.CanalServerWithEmbedded;
 import com.alibaba.otter.canal.server.embedded.CanalServerWithEmbedded;
@@ -146,7 +145,7 @@ public class CanalMQStarter {
                     try {
                     try {
                         int size = message.isRaw() ? message.getRawEntries().size() : message.getEntries().size();
                         int size = message.isRaw() ? message.getRawEntries().size() : message.getEntries().size();
                         if (batchId != -1 && size != 0) {
                         if (batchId != -1 && size != 0) {
-                            canalMQProducer.send(destination, message, new CanalKafkaProducer.Callback() {
+                            canalMQProducer.send(destination, message, new CanalMQProducer.Callback() {
 
 
                                 @Override
                                 @Override
                                 public void commit() {
                                 public void commit() {