mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
Merge pull request #3683 from duhenglucky/develop_jdk
[ISSUE #3684] change client jdk version to 1.6
This commit is contained in:
@@ -27,6 +27,11 @@
|
||||
<artifactId>rocketmq-client</artifactId>
|
||||
<name>rocketmq-client ${project.version}</name>
|
||||
|
||||
<properties>
|
||||
<maven.compiler.source>1.6</maven.compiler.source>
|
||||
<maven.compiler.target>1.6</maven.compiler.target>
|
||||
</properties>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>${project.groupId}</groupId>
|
||||
|
||||
@@ -397,7 +397,8 @@ public class MQClientAPIImpl {
|
||||
|
||||
}
|
||||
|
||||
public AclConfig getBrokerClusterConfig(final String addr, final long timeoutMillis) throws RemotingCommandException, InterruptedException, RemotingTimeoutException,
|
||||
public AclConfig getBrokerClusterConfig(final String addr,
|
||||
final long timeoutMillis) throws RemotingCommandException, InterruptedException, RemotingTimeoutException,
|
||||
RemotingSendRequestException, RemotingConnectException, MQBrokerException {
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_BROKER_CLUSTER_ACL_CONFIG, null);
|
||||
|
||||
@@ -407,7 +408,7 @@ public class MQClientAPIImpl {
|
||||
case ResponseCode.SUCCESS: {
|
||||
if (response.getBody() != null) {
|
||||
GetBrokerClusterAclConfigResponseBody body =
|
||||
GetBrokerClusterAclConfigResponseBody.decode(response.getBody(), GetBrokerClusterAclConfigResponseBody.class);
|
||||
GetBrokerClusterAclConfigResponseBody.decode(response.getBody(), GetBrokerClusterAclConfigResponseBody.class);
|
||||
AclConfig aclConfig = new AclConfig();
|
||||
aclConfig.setGlobalWhiteAddrs(body.getGlobalWhiteAddrs());
|
||||
aclConfig.setPlainAccessConfigs(body.getPlainAccessConfigs());
|
||||
@@ -505,7 +506,7 @@ public class MQClientAPIImpl {
|
||||
) throws RemotingException, MQBrokerException, InterruptedException {
|
||||
RemotingCommand response = this.remotingClient.invokeSync(addr, request, timeoutMillis);
|
||||
assert response != null;
|
||||
return this.processSendResponse(brokerName, msg, response,addr);
|
||||
return this.processSendResponse(brokerName, msg, response, addr);
|
||||
}
|
||||
|
||||
private void sendMessageAsync(
|
||||
@@ -619,7 +620,10 @@ public class MQClientAPIImpl {
|
||||
request.setOpaque(RemotingCommand.createNewRequestId());
|
||||
sendMessageAsync(addr, retryBrokerName, msg, timeoutMillis, request, sendCallback, topicPublishInfo, instance,
|
||||
timesTotal, curTimes, context, producer);
|
||||
} catch (InterruptedException | RemotingTooMuchRequestException e1) {
|
||||
} catch (InterruptedException e1) {
|
||||
onExceptionImpl(retryBrokerName, msg, timeoutMillis, request, sendCallback, topicPublishInfo, instance, timesTotal, curTimes, e1,
|
||||
context, false, producer);
|
||||
} catch (RemotingTooMuchRequestException e1) {
|
||||
onExceptionImpl(retryBrokerName, msg, timeoutMillis, request, sendCallback, topicPublishInfo, instance, timesTotal, curTimes, e1,
|
||||
context, false, producer);
|
||||
} catch (RemotingException e1) {
|
||||
@@ -671,7 +675,7 @@ public class MQClientAPIImpl {
|
||||
}
|
||||
|
||||
SendMessageResponseHeader responseHeader =
|
||||
(SendMessageResponseHeader) response.decodeCommandCustomHeader(SendMessageResponseHeader.class);
|
||||
(SendMessageResponseHeader) response.decodeCommandCustomHeader(SendMessageResponseHeader.class);
|
||||
|
||||
//If namespace not null , reset Topic without namespace.
|
||||
String topic = msg.getTopic();
|
||||
@@ -690,8 +694,8 @@ public class MQClientAPIImpl {
|
||||
uniqMsgId = sb.toString();
|
||||
}
|
||||
SendResult sendResult = new SendResult(sendStatus,
|
||||
uniqMsgId,
|
||||
responseHeader.getMsgId(), messageQueue, responseHeader.getQueueOffset());
|
||||
uniqMsgId,
|
||||
responseHeader.getMsgId(), messageQueue, responseHeader.getQueueOffset());
|
||||
sendResult.setTransactionId(responseHeader.getTransactionId());
|
||||
String regionId = response.getExtFields().get(MessageConst.PROPERTY_MSG_REGION);
|
||||
String traceOn = response.getExtFields().get(MessageConst.PROPERTY_TRACE_SWITCH);
|
||||
@@ -1432,8 +1436,8 @@ public class MQClientAPIImpl {
|
||||
}
|
||||
|
||||
public int addWritePermOfBroker(final String nameSrvAddr, String brokerName, final long timeoutMillis)
|
||||
throws RemotingCommandException,
|
||||
RemotingConnectException, RemotingSendRequestException, RemotingTimeoutException, InterruptedException, MQClientException {
|
||||
throws RemotingCommandException,
|
||||
RemotingConnectException, RemotingSendRequestException, RemotingTimeoutException, InterruptedException, MQClientException {
|
||||
AddWritePermOfBrokerRequestHeader requestHeader = new AddWritePermOfBrokerRequestHeader();
|
||||
requestHeader.setBrokerName(brokerName);
|
||||
|
||||
@@ -1444,7 +1448,7 @@ public class MQClientAPIImpl {
|
||||
switch (response.getCode()) {
|
||||
case ResponseCode.SUCCESS: {
|
||||
AddWritePermOfBrokerResponseHeader responseHeader =
|
||||
(AddWritePermOfBrokerResponseHeader) response.decodeCommandCustomHeader(AddWritePermOfBrokerResponseHeader.class);
|
||||
(AddWritePermOfBrokerResponseHeader) response.decodeCommandCustomHeader(AddWritePermOfBrokerResponseHeader.class);
|
||||
return responseHeader.getAddTopicCount();
|
||||
}
|
||||
default:
|
||||
@@ -1492,7 +1496,8 @@ public class MQClientAPIImpl {
|
||||
throw new MQClientException(response.getCode(), response.getRemark());
|
||||
}
|
||||
|
||||
public void deleteSubscriptionGroup(final String addr, final String groupName, final boolean removeOffset, final long timeoutMillis)
|
||||
public void deleteSubscriptionGroup(final String addr, final String groupName, final boolean removeOffset,
|
||||
final long timeoutMillis)
|
||||
throws RemotingException, InterruptedException, MQClientException {
|
||||
DeleteSubscriptionGroupRequestHeader requestHeader = new DeleteSubscriptionGroupRequestHeader();
|
||||
requestHeader.setGroupName(groupName);
|
||||
|
||||
+1
-1
@@ -147,7 +147,7 @@ public class DefaultLitePullConsumerImpl implements MQConsumerInner {
|
||||
|
||||
private final MessageQueueLock messageQueueLock = new MessageQueueLock();
|
||||
|
||||
private final ArrayList<ConsumeMessageHook> consumeMessageHookList = new ArrayList<>();
|
||||
private final ArrayList<ConsumeMessageHook> consumeMessageHookList = new ArrayList<ConsumeMessageHook>();
|
||||
|
||||
public DefaultLitePullConsumerImpl(final DefaultLitePullConsumer defaultLitePullConsumer, final RPCHook rpcHook) {
|
||||
this.defaultLitePullConsumer = defaultLitePullConsumer;
|
||||
|
||||
@@ -897,8 +897,12 @@ public class MQClientInstance {
|
||||
try {
|
||||
this.mQClientAPIImpl.unregisterClient(addr, this.clientId, producerGroup, consumerGroup, clientConfig.getMqClientApiTimeout());
|
||||
log.info("unregister client[Producer: {} Consumer: {}] from broker[{} {} {}] success", producerGroup, consumerGroup, brokerName, entry1.getKey(), addr);
|
||||
} catch (RemotingException | InterruptedException | MQBrokerException e) {
|
||||
log.error("unregister client exception from broker: " + addr, e);
|
||||
} catch (RemotingException e) {
|
||||
log.warn("unregister client RemotingException from broker: {}, {}", addr, e.getMessage());
|
||||
} catch (InterruptedException e) {
|
||||
log.warn("unregister client InterruptedException from broker: {}, {}", addr, e.getMessage());
|
||||
} catch (MQBrokerException e) {
|
||||
log.warn("unregister client MQBrokerException from broker: {}, {}", addr, e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+27
-6
@@ -490,7 +490,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
*
|
||||
* @param msg
|
||||
* @param sendCallback
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @throws RejectedExecutionException
|
||||
*/
|
||||
@Deprecated
|
||||
@@ -597,7 +597,14 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
default:
|
||||
break;
|
||||
}
|
||||
} catch (RemotingException | MQClientException e) {
|
||||
} catch (RemotingException e) {
|
||||
endTimestamp = System.currentTimeMillis();
|
||||
this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, true);
|
||||
log.warn(String.format("sendKernelImpl exception, resend at once, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
|
||||
log.warn(msg.toString());
|
||||
exception = e;
|
||||
continue;
|
||||
} catch (MQClientException e) {
|
||||
endTimestamp = System.currentTimeMillis();
|
||||
this.updateFaultItem(mq.getBrokerName(), endTimestamp - beginTimestampPrev, true);
|
||||
log.warn(String.format("sendKernelImpl exception, resend at once, InvokeID: %s, RT: %sms, Broker: %s", invokeID, endTimestamp - beginTimestampPrev, mq), e);
|
||||
@@ -856,7 +863,19 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
}
|
||||
|
||||
return sendResult;
|
||||
} catch (RemotingException | MQBrokerException | InterruptedException e) {
|
||||
} catch (RemotingException e) {
|
||||
if (this.hasSendMessageHook()) {
|
||||
context.setException(e);
|
||||
this.executeSendMessageHookAfter(context);
|
||||
}
|
||||
throw e;
|
||||
} catch (MQBrokerException e) {
|
||||
if (this.hasSendMessageHook()) {
|
||||
context.setException(e);
|
||||
this.executeSendMessageHookAfter(context);
|
||||
}
|
||||
throw e;
|
||||
} catch (InterruptedException e) {
|
||||
if (this.hasSendMessageHook()) {
|
||||
context.setException(e);
|
||||
this.executeSendMessageHookAfter(context);
|
||||
@@ -969,6 +988,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
executeEndTransactionHook(context);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* DEFAULT ONEWAY -------------------------------------------------------
|
||||
*/
|
||||
@@ -1021,7 +1041,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
* @param msg
|
||||
* @param mq
|
||||
* @param sendCallback
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @throws MQClientException
|
||||
* @throws RemotingException
|
||||
* @throws InterruptedException
|
||||
@@ -1151,7 +1171,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
* @param selector
|
||||
* @param arg
|
||||
* @param sendCallback
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @param timeout the <code>sendCallback</code> will be invoked at most time
|
||||
* @throws MQClientException
|
||||
* @throws RemotingException
|
||||
* @throws InterruptedException
|
||||
@@ -1491,7 +1511,8 @@ public class DefaultMQProducerImpl implements MQProducerInner {
|
||||
}
|
||||
}
|
||||
|
||||
private Message waitResponse(Message msg, long timeout, RequestResponseFuture requestResponseFuture, long cost) throws InterruptedException, RequestTimeoutException, MQClientException {
|
||||
private Message waitResponse(Message msg, long timeout, RequestResponseFuture requestResponseFuture,
|
||||
long cost) throws InterruptedException, RequestTimeoutException, MQClientException {
|
||||
Message responseMessage = requestResponseFuture.waitResponseMessage(timeout - cost);
|
||||
if (responseMessage == null) {
|
||||
if (requestResponseFuture.isSendRequestOk()) {
|
||||
|
||||
@@ -39,7 +39,7 @@ public class RequestFutureHolder {
|
||||
private static InternalLogger log = ClientLogger.getLog();
|
||||
private static final RequestFutureHolder INSTANCE = new RequestFutureHolder();
|
||||
private ConcurrentHashMap<String, RequestResponseFuture> requestFutureTable = new ConcurrentHashMap<String, RequestResponseFuture>();
|
||||
private final Set<DefaultMQProducerImpl> producerSet = new HashSet<>();
|
||||
private final Set<DefaultMQProducerImpl> producerSet = new HashSet<DefaultMQProducerImpl>();
|
||||
private ScheduledExecutorService scheduledExecutorService = null;
|
||||
|
||||
public ConcurrentHashMap<String, RequestResponseFuture> getRequestFutureTable() {
|
||||
|
||||
+1
-1
@@ -51,7 +51,7 @@ public class ConsumeMessageOpenTracingHookImpl implements ConsumeMessageHook {
|
||||
if (context == null || context.getMsgList() == null || context.getMsgList().isEmpty()) {
|
||||
return;
|
||||
}
|
||||
List<Span> spanList = new ArrayList<>();
|
||||
List<Span> spanList = new ArrayList<Span>();
|
||||
for (MessageExt msg : context.getMsgList()) {
|
||||
if (msg == null) {
|
||||
continue;
|
||||
|
||||
+11
-8
@@ -23,6 +23,7 @@ import java.net.InetSocketAddress;
|
||||
import java.util.Collections;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import org.apache.commons.lang3.reflect.FieldUtils;
|
||||
@@ -95,7 +96,9 @@ public class DefaultLitePullConsumerTest {
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ConcurrentMap<String, MQClientInstance> factoryTable = (ConcurrentMap<String, MQClientInstance>) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true);
|
||||
factoryTable.forEach((s, instance) -> instance.shutdown());
|
||||
for (Map.Entry<String, MQClientInstance> entry : factoryTable.entrySet()) {
|
||||
entry.getValue().shutdown();
|
||||
}
|
||||
factoryTable.clear();
|
||||
|
||||
Field field = MQClientInstance.class.getDeclaredField("rebalanceService");
|
||||
@@ -481,7 +484,7 @@ public class DefaultLitePullConsumerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testComputePullFromWhereReturnedNotFound() throws Exception{
|
||||
public void testComputePullFromWhereReturnedNotFound() throws Exception {
|
||||
DefaultLitePullConsumer defaultLitePullConsumer = createStartLitePullConsumer();
|
||||
defaultLitePullConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
|
||||
MessageQueue messageQueue = createMessageQueue();
|
||||
@@ -491,7 +494,7 @@ public class DefaultLitePullConsumerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testComputePullFromWhereReturned() throws Exception{
|
||||
public void testComputePullFromWhereReturned() throws Exception {
|
||||
DefaultLitePullConsumer defaultLitePullConsumer = createStartLitePullConsumer();
|
||||
defaultLitePullConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
|
||||
MessageQueue messageQueue = createMessageQueue();
|
||||
@@ -500,9 +503,8 @@ public class DefaultLitePullConsumerTest {
|
||||
assertThat(offset).isEqualTo(100);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testComputePullFromLast() throws Exception{
|
||||
public void testComputePullFromLast() throws Exception {
|
||||
DefaultLitePullConsumer defaultLitePullConsumer = createStartLitePullConsumer();
|
||||
defaultLitePullConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
|
||||
MessageQueue messageQueue = createMessageQueue();
|
||||
@@ -513,13 +515,13 @@ public class DefaultLitePullConsumerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testComputePullByTimeStamp() throws Exception{
|
||||
public void testComputePullByTimeStamp() throws Exception {
|
||||
DefaultLitePullConsumer defaultLitePullConsumer = createStartLitePullConsumer();
|
||||
defaultLitePullConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_TIMESTAMP);
|
||||
defaultLitePullConsumer.setConsumeTimestamp("20191024171201");
|
||||
MessageQueue messageQueue = createMessageQueue();
|
||||
when(offsetStore.readOffset(any(MessageQueue.class), any(ReadOffsetType.class))).thenReturn(-1L);
|
||||
when(mQClientFactory.getMQAdminImpl().searchOffset(any(MessageQueue.class),anyLong())).thenReturn(100L);
|
||||
when(mQClientFactory.getMQAdminImpl().searchOffset(any(MessageQueue.class), anyLong())).thenReturn(100L);
|
||||
long offset = rebalanceImpl.computePullFromWhere(messageQueue);
|
||||
assertThat(offset).isEqualTo(100);
|
||||
}
|
||||
@@ -660,7 +662,8 @@ public class DefaultLitePullConsumerTest {
|
||||
return new PullResultExt(pullStatus, requestHeader.getQueueOffset() + messageExtList.size(), 123, 2048, messageExtList, 0, outputStream.toByteArray());
|
||||
}
|
||||
|
||||
private static void suppressUpdateTopicRouteInfoFromNameServer(DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException {
|
||||
private static void suppressUpdateTopicRouteInfoFromNameServer(
|
||||
DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException {
|
||||
DefaultLitePullConsumerImpl defaultLitePullConsumerImpl = (DefaultLitePullConsumerImpl) FieldUtils.readDeclaredField(litePullConsumer, "defaultLitePullConsumerImpl", true);
|
||||
if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) {
|
||||
litePullConsumer.changeInstanceNameToPID();
|
||||
|
||||
+13
-9
@@ -21,6 +21,7 @@ import java.net.InetSocketAddress;
|
||||
import java.util.Collections;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
@@ -94,7 +95,9 @@ public class DefaultMQPushConsumerTest {
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ConcurrentMap<String, MQClientInstance> factoryTable = (ConcurrentMap<String, MQClientInstance>) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true);
|
||||
factoryTable.forEach((s, instance) -> instance.shutdown());
|
||||
for (Map.Entry<String, MQClientInstance> entry : factoryTable.entrySet()) {
|
||||
entry.getValue().shutdown();
|
||||
}
|
||||
factoryTable.clear();
|
||||
|
||||
when(mQClientAPIImpl.pullMessage(anyString(), any(PullMessageRequestHeader.class),
|
||||
@@ -117,7 +120,6 @@ public class DefaultMQPushConsumerTest {
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
consumerGroup = "FooBarGroup" + System.currentTimeMillis();
|
||||
pushConsumer = new DefaultMQPushConsumer(consumerGroup);
|
||||
pushConsumer.setNamesrvAddr("127.0.0.1:9876");
|
||||
@@ -169,7 +171,7 @@ public class DefaultMQPushConsumerTest {
|
||||
@Test
|
||||
public void testPullMessage_Success() throws InterruptedException, RemotingException, MQBrokerException {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<>();
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<MessageExt>();
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().setConsumeMessageService(new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(), new MessageListenerConcurrently() {
|
||||
@Override
|
||||
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
|
||||
@@ -192,7 +194,7 @@ public class DefaultMQPushConsumerTest {
|
||||
@Test
|
||||
public void testPullMessage_SuccessWithOrderlyService() throws Exception {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<>();
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<MessageExt>();
|
||||
|
||||
MessageListenerOrderly listenerOrderly = new MessageListenerOrderly() {
|
||||
@Override
|
||||
@@ -329,11 +331,13 @@ public class DefaultMQPushConsumerTest {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final MessageExt[] messageExts = new MessageExt[1];
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().setConsumeMessageService(
|
||||
new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(),
|
||||
(msgs, context) -> {
|
||||
messageExts[0] = msgs.get(0);
|
||||
return null;
|
||||
}));
|
||||
new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(), new MessageListenerConcurrently() {
|
||||
@Override public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
|
||||
ConsumeConcurrentlyContext context) {
|
||||
messageExts[0] = msgs.get(0);
|
||||
return null;
|
||||
}
|
||||
}));
|
||||
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().setConsumeOrderly(true);
|
||||
PullMessageService pullMessageService = mQClientFactory.getPullMessageService();
|
||||
|
||||
+1
-1
@@ -149,7 +149,7 @@ public class ConsumeMessageConcurrentlyServiceTest {
|
||||
@Test
|
||||
public void testPullMessage_ConsumeSuccess() throws InterruptedException, RemotingException, MQBrokerException, NoSuchFieldException,Exception {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<>();
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<MessageExt>();
|
||||
|
||||
ConsumeMessageConcurrentlyService normalServie = new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(), new MessageListenerConcurrently() {
|
||||
@Override
|
||||
|
||||
@@ -257,7 +257,7 @@ public class DefaultMQProducerTest {
|
||||
}
|
||||
};
|
||||
|
||||
List<Message> msgs = new ArrayList<>();
|
||||
List<Message> msgs = new ArrayList<Message>();
|
||||
for (int i = 0; i < 5; i++) {
|
||||
Message message = new Message();
|
||||
message.setTopic("test");
|
||||
|
||||
+19
-11
@@ -25,7 +25,9 @@ import java.net.InetSocketAddress;
|
||||
import java.util.Collections;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -97,12 +99,14 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
private PullAPIWrapper pullAPIWrapper;
|
||||
private RebalancePushImpl rebalancePushImpl;
|
||||
private DefaultMQPushConsumer pushConsumer;
|
||||
private MockTracer tracer = new MockTracer();
|
||||
private final MockTracer tracer = new MockTracer();
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ConcurrentMap<String, MQClientInstance> factoryTable = (ConcurrentMap<String, MQClientInstance>) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true);
|
||||
factoryTable.forEach((s, instance) -> instance.shutdown());
|
||||
for (Map.Entry<String, MQClientInstance> entry : factoryTable.entrySet()) {
|
||||
entry.getValue().shutdown();
|
||||
}
|
||||
factoryTable.clear();
|
||||
|
||||
when(mQClientAPIImpl.pullMessage(anyString(), any(PullMessageRequestHeader.class),
|
||||
@@ -115,7 +119,7 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
messageClientExt.setTopic(topic);
|
||||
messageClientExt.setQueueId(0);
|
||||
messageClientExt.setMsgId("123");
|
||||
messageClientExt.setBody(new byte[]{'a'});
|
||||
messageClientExt.setBody(new byte[] {'a'});
|
||||
messageClientExt.setOffsetMsgId("234");
|
||||
messageClientExt.setBornHost(new InetSocketAddress(8080));
|
||||
messageClientExt.setStoreHost(new InetSocketAddress(8080));
|
||||
@@ -125,11 +129,10 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
consumerGroup = "FooBarGroup" + System.currentTimeMillis();
|
||||
pushConsumer = new DefaultMQPushConsumer(consumerGroup);
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().registerConsumeMessageHook(
|
||||
new ConsumeMessageOpenTracingHookImpl(tracer));
|
||||
new ConsumeMessageOpenTracingHookImpl(tracer));
|
||||
pushConsumer.setNamesrvAddr("127.0.0.1:9876");
|
||||
pushConsumer.setPullInterval(60 * 1000);
|
||||
|
||||
@@ -140,7 +143,7 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
pushConsumer.registerMessageListener(new MessageListenerConcurrently() {
|
||||
@Override
|
||||
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
|
||||
ConsumeConcurrentlyContext context) {
|
||||
ConsumeConcurrentlyContext context) {
|
||||
return null;
|
||||
}
|
||||
});
|
||||
@@ -173,11 +176,11 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
@Test
|
||||
public void testPullMessage_WithTrace_Success() throws InterruptedException, RemotingException, MQBrokerException, MQClientException {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<>();
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<MessageExt>();
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().setConsumeMessageService(new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(), new MessageListenerConcurrently() {
|
||||
@Override
|
||||
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
|
||||
ConsumeConcurrentlyContext context) {
|
||||
ConsumeConcurrentlyContext context) {
|
||||
messageAtomic.set(msgs.get(0));
|
||||
countDownLatch.countDown();
|
||||
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
|
||||
@@ -190,10 +193,15 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
MessageExt msg = messageAtomic.get();
|
||||
assertThat(msg).isNotNull();
|
||||
assertThat(msg.getTopic()).isEqualTo(topic);
|
||||
assertThat(msg.getBody()).isEqualTo(new byte[]{'a'});
|
||||
assertThat(msg.getBody()).isEqualTo(new byte[] {'a'});
|
||||
|
||||
// wait until consumeMessageAfter hook of tracer is done surely.
|
||||
waitAtMost(1, TimeUnit.SECONDS).until(() -> tracer.finishedSpans().size() == 1);
|
||||
waitAtMost(1, TimeUnit.SECONDS).until(new Callable() {
|
||||
@Override public Object call() throws Exception {
|
||||
return tracer.finishedSpans().size() == 1;
|
||||
}
|
||||
});
|
||||
|
||||
MockSpan span = tracer.finishedSpans().get(0);
|
||||
assertThat(span.tags().get(Tags.MESSAGE_BUS_DESTINATION.getKey())).isEqualTo(topic);
|
||||
assertThat(span.tags().get(Tags.SPAN_KIND.getKey())).isEqualTo(Tags.SPAN_KIND_CONSUMER);
|
||||
@@ -219,7 +227,7 @@ public class DefaultMQConsumerWithOpenTracingTest {
|
||||
}
|
||||
|
||||
private PullResultExt createPullResult(PullMessageRequestHeader requestHeader, PullStatus pullStatus,
|
||||
List<MessageExt> messageExtList) throws Exception {
|
||||
List<MessageExt> messageExtList) throws Exception {
|
||||
ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
|
||||
for (MessageExt messageExt : messageExtList) {
|
||||
outputStream.write(MessageDecoder.encode(messageExt, false));
|
||||
|
||||
+5
-2
@@ -25,6 +25,7 @@ import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
@@ -116,7 +117,9 @@ public class DefaultMQConsumerWithTraceTest {
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ConcurrentMap<String, MQClientInstance> factoryTable = (ConcurrentMap<String, MQClientInstance>) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true);
|
||||
factoryTable.forEach((s, instance) -> instance.shutdown());
|
||||
for (Map.Entry<String, MQClientInstance> entry : factoryTable.entrySet()) {
|
||||
entry.getValue().shutdown();
|
||||
}
|
||||
factoryTable.clear();
|
||||
|
||||
consumerGroup = "FooBarGroup" + System.currentTimeMillis();
|
||||
@@ -217,7 +220,7 @@ public class DefaultMQConsumerWithTraceTest {
|
||||
traceProducer.getDefaultMQProducerImpl().getmQClientFactory().registerProducer(producerGroupTraceTemp, traceProducer.getDefaultMQProducerImpl());
|
||||
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<>();
|
||||
final AtomicReference<MessageExt> messageAtomic = new AtomicReference<MessageExt>();
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().setConsumeMessageService(new ConsumeMessageConcurrentlyService(pushConsumer.getDefaultMQPushConsumerImpl(), new MessageListenerConcurrently() {
|
||||
@Override
|
||||
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs,
|
||||
|
||||
+22
-14
@@ -17,6 +17,12 @@
|
||||
|
||||
package org.apache.rocketmq.client.trace;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.apache.rocketmq.client.ClientConfig;
|
||||
import org.apache.rocketmq.client.exception.MQBrokerException;
|
||||
import org.apache.rocketmq.client.exception.MQClientException;
|
||||
@@ -52,17 +58,16 @@ import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.Spy;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.junit.MockitoJUnitRunner;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import org.mockito.stubbing.Answer;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.ArgumentMatchers.*;
|
||||
import static org.mockito.ArgumentMatchers.any;
|
||||
import static org.mockito.ArgumentMatchers.anyInt;
|
||||
import static org.mockito.ArgumentMatchers.anyLong;
|
||||
import static org.mockito.ArgumentMatchers.anyString;
|
||||
import static org.mockito.ArgumentMatchers.nullable;
|
||||
import static org.mockito.Mockito.doAnswer;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
@@ -127,7 +132,7 @@ public class TransactionMQProducerWithTraceTest {
|
||||
|
||||
Field fieldHooks = DefaultMQProducerImpl.class.getDeclaredField("endTransactionHookList");
|
||||
fieldHooks.setAccessible(true);
|
||||
List<EndTransactionHook>hooks = new ArrayList<>();
|
||||
List<EndTransactionHook> hooks = new ArrayList<EndTransactionHook>();
|
||||
hooks.add(endTransactionHook);
|
||||
fieldHooks.set(producer.getDefaultMQProducerImpl(), hooks);
|
||||
|
||||
@@ -143,11 +148,14 @@ public class TransactionMQProducerWithTraceTest {
|
||||
public void testSendMessageSync_WithTrace_Success() throws RemotingException, InterruptedException, MQBrokerException, MQClientException {
|
||||
traceProducer.getDefaultMQProducerImpl().getmQClientFactory().registerProducer(producerGroupTraceTemp, traceProducer.getDefaultMQProducerImpl());
|
||||
when(mQClientAPIImpl.getTopicRouteInfoFromNameServer(anyString(), anyLong())).thenReturn(createTopicRoute());
|
||||
AtomicReference<EndTransactionContext> context = new AtomicReference<>();
|
||||
doAnswer(mock -> {
|
||||
context.set(mock.getArgument(0));
|
||||
return null;
|
||||
}).when(endTransactionHook).endTransaction(any());
|
||||
final AtomicReference<EndTransactionContext> context = new AtomicReference<EndTransactionContext>();
|
||||
doAnswer(new Answer() {
|
||||
@Override public Object answer(InvocationOnMock mock) throws Throwable {
|
||||
context.set((EndTransactionContext) mock.getArgument(0));
|
||||
return null;
|
||||
}
|
||||
|
||||
}).when(endTransactionHook).endTransaction(any(EndTransactionContext.class));
|
||||
producer.sendMessageInTransaction(message, null);
|
||||
|
||||
EndTransactionContext ctx = context.get();
|
||||
|
||||
Reference in New Issue
Block a user