Skip to content

Commit 54fda9b

Browse files
committed
upgrade rocketmq client version to 5.1.0
1 parent f697554 commit 54fda9b

File tree

8 files changed

+8
-16
lines changed

8 files changed

+8
-16
lines changed

rocketmq-connect-runtime/src/main/java/org/apache/rocketmq/connect/runtime/stats/ConnectStatsService.java

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,6 @@
2020
import org.apache.commons.lang3.StringUtils;
2121
import org.apache.rocketmq.common.ServiceThread;
2222
import org.apache.rocketmq.connect.runtime.common.LoggerName;
23-
import org.apache.rocketmq.logging.InternalLogger;
24-
import org.apache.rocketmq.logging.InternalLoggerFactory;
2523

2624
import java.text.MessageFormat;
2725
import java.util.HashMap;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/connectorwrapper/NameServerMocker.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,12 @@
1717
package org.apache.rocketmq.connect.runtime.connectorwrapper;
1818

1919
import org.apache.rocketmq.common.MixAll;
20-
import org.apache.rocketmq.common.protocol.route.BrokerData;
21-
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
2220

2321
import java.util.ArrayList;
2422
import java.util.HashMap;
2523
import java.util.List;
24+
import org.apache.rocketmq.remoting.protocol.route.BrokerData;
25+
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
2626

2727
/**
2828
* tools class

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/service/ClusterManagementServiceImplTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,13 @@
2121
import io.netty.channel.ChannelHandlerContext;
2222
import java.nio.charset.StandardCharsets;
2323
import java.util.List;
24-
import org.apache.rocketmq.common.protocol.RequestCode;
25-
import org.apache.rocketmq.common.protocol.header.NotifyConsumerIdsChangedRequestHeader;
2624
import org.apache.rocketmq.connect.runtime.config.WorkerConfig;
2725
import org.apache.rocketmq.connect.runtime.connectorwrapper.NameServerMocker;
2826
import org.apache.rocketmq.connect.runtime.connectorwrapper.ServerResponseMocker;
2927
import org.apache.rocketmq.remoting.RemotingClient;
3028
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
29+
import org.apache.rocketmq.remoting.protocol.RequestCode;
30+
import org.apache.rocketmq.remoting.protocol.header.NotifyConsumerIdsChangedRequestHeader;
3131
import org.junit.After;
3232
import org.junit.Assert;
3333
import org.junit.Before;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/service/ConfigManagementServiceImplTest.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -29,12 +29,9 @@
2929
import java.util.Set;
3030
import java.util.UUID;
3131

32-
import com.google.common.collect.Maps;
3332
import org.apache.rocketmq.client.consumer.DefaultLitePullConsumer;
34-
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
3533
import org.apache.rocketmq.client.producer.DefaultMQProducer;
3634
import org.apache.rocketmq.client.producer.SendCallback;
37-
import org.apache.rocketmq.common.admin.TopicOffset;
3835
import org.apache.rocketmq.common.message.Message;
3936
import org.apache.rocketmq.common.message.MessageQueue;
4037
import org.apache.rocketmq.connect.runtime.common.ConnAndTaskConfigs;
@@ -53,6 +50,7 @@
5350
import org.apache.rocketmq.connect.runtime.utils.datasync.BrokerBasedLog;
5451
import org.apache.rocketmq.connect.runtime.utils.datasync.DataSynchronizer;
5552
import org.apache.rocketmq.connect.runtime.utils.datasync.DataSynchronizerCallback;
53+
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
5654
import org.junit.After;
5755
import org.junit.Assert;
5856
import org.junit.Before;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/service/DefaultConnectorContextTest.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,6 @@
4040
import org.apache.rocketmq.client.impl.consumer.RebalanceImpl;
4141
import org.apache.rocketmq.client.impl.factory.MQClientInstance;
4242
import org.apache.rocketmq.client.producer.DefaultMQProducer;
43-
import org.apache.rocketmq.common.admin.TopicOffset;
4443
import org.apache.rocketmq.common.message.MessageQueue;
4544
import org.apache.rocketmq.connect.runtime.common.ConnectKeyValue;
4645
import org.apache.rocketmq.connect.runtime.config.WorkerConfig;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/service/PositionManagementServiceImplTest.java

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,16 +19,12 @@
1919

2020
import com.google.common.collect.Lists;
2121
import io.netty.util.internal.ConcurrentSet;
22-
import io.openmessaging.Future;
2322
import io.openmessaging.connector.api.data.RecordOffset;
24-
import io.openmessaging.producer.SendResult;
2523
import org.apache.rocketmq.client.consumer.DefaultLitePullConsumer;
2624
import org.apache.rocketmq.client.producer.DefaultMQProducer;
2725
import org.apache.rocketmq.client.producer.SendCallback;
28-
import org.apache.rocketmq.common.admin.TopicOffset;
2926
import org.apache.rocketmq.common.message.Message;
3027
import org.apache.rocketmq.common.message.MessageQueue;
31-
import org.apache.rocketmq.connect.runtime.common.ConnAndTaskConfigs;
3228
import org.apache.rocketmq.connect.runtime.config.WorkerConfig;
3329
import org.apache.rocketmq.connect.runtime.converter.record.json.JsonConverter;
3430
import org.apache.rocketmq.connect.runtime.connectorwrapper.NameServerMocker;
@@ -40,6 +36,7 @@
4036
import org.apache.rocketmq.connect.runtime.utils.TestUtils;
4137
import org.apache.rocketmq.connect.runtime.utils.datasync.BrokerBasedLog;
4238
import org.apache.rocketmq.connect.runtime.utils.datasync.DataSynchronizerCallback;
39+
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
4340
import org.assertj.core.util.Maps;
4441
import org.junit.After;
4542
import org.junit.Before;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/store/PositionStorageReaderImplTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
import io.openmessaging.connector.api.data.RecordPartition;
2525
import org.apache.rocketmq.client.consumer.DefaultLitePullConsumer;
2626
import org.apache.rocketmq.client.producer.DefaultMQProducer;
27-
import org.apache.rocketmq.common.admin.TopicOffset;
2827
import org.apache.rocketmq.common.message.MessageQueue;
2928
import org.apache.rocketmq.connect.runtime.config.WorkerConfig;
3029
import org.apache.rocketmq.connect.runtime.connectorwrapper.NameServerMocker;
@@ -33,6 +32,7 @@
3332
import org.apache.rocketmq.connect.runtime.service.PositionManagementService;
3433
import org.apache.rocketmq.connect.runtime.service.local.LocalPositionManagementServiceImpl;
3534
import org.apache.rocketmq.connect.runtime.utils.ConnectUtil;
35+
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
3636
import org.assertj.core.util.Maps;
3737
import org.junit.After;
3838
import org.junit.Assert;

rocketmq-connect-runtime/src/test/java/org/apache/rocketmq/connect/runtime/utils/datasync/BrokerBasedLogTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,14 +23,14 @@
2323
import org.apache.rocketmq.client.producer.DefaultMQProducer;
2424
import org.apache.rocketmq.client.producer.SendCallback;
2525
import org.apache.rocketmq.client.producer.selector.SelectMessageQueueByHash;
26-
import org.apache.rocketmq.common.admin.TopicOffset;
2726
import org.apache.rocketmq.common.message.Message;
2827
import org.apache.rocketmq.common.message.MessageQueue;
2928
import org.apache.rocketmq.connect.runtime.config.WorkerConfig;
3029
import org.apache.rocketmq.connect.runtime.serialization.Serde;
3130
import org.apache.rocketmq.connect.runtime.serialization.Serializer;
3231
import org.apache.rocketmq.connect.runtime.utils.ConnectUtil;
3332
import org.apache.rocketmq.remoting.exception.RemotingException;
33+
import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
3434
import org.junit.After;
3535
import org.junit.Before;
3636
import org.junit.Test;

0 commit comments

Comments
 (0)