|
@@ -12,7 +12,6 @@ import org.apache.rocketmq.client.consumer.listener.ConsumeOrderlyStatus;
|
|
|
import org.apache.rocketmq.client.consumer.listener.MessageListenerOrderly;
|
|
|
import org.apache.rocketmq.client.consumer.rebalance.AllocateMessageQueueAveragely;
|
|
|
import org.apache.rocketmq.client.exception.MQClientException;
|
|
|
-import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
|
|
import org.apache.rocketmq.common.message.MessageExt;
|
|
|
import org.apache.rocketmq.remoting.RPCHook;
|
|
|
import org.slf4j.Logger;
|
|
@@ -84,7 +83,6 @@ public class RocketMQCanalConnector implements CanalMQConnector {
|
|
|
}
|
|
|
rocketMQConsumer = new DefaultMQPushConsumer(groupName, rpcHook, new AllocateMessageQueueAveragely());
|
|
|
rocketMQConsumer.setVipChannelEnabled(false);
|
|
|
- rocketMQConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
|
|
|
if (!StringUtils.isBlank(nameServer)) {
|
|
|
rocketMQConsumer.setNamesrvAddr(nameServer);
|
|
|
}
|