mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
This commit is contained in:
+4
-2
@@ -174,8 +174,10 @@ public class NotificationProcessor implements NettyRequestProcessor {
|
||||
}
|
||||
|
||||
private boolean hasMsgFromQueue(boolean isRetry, NotificationRequestHeader requestHeader, int queueId) {
|
||||
if (this.brokerController.getConsumerOrderInfoManager().checkBlock(null, requestHeader.getTopic(), requestHeader.getConsumerGroup(), queueId, 0)) {
|
||||
return false;
|
||||
if (Boolean.TRUE.equals(requestHeader.getOrder())) {
|
||||
if (this.brokerController.getConsumerOrderInfoManager().checkBlock(requestHeader.getAttemptId(), requestHeader.getTopic(), requestHeader.getConsumerGroup(), queueId, 0)) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
String topic = isRetry ? KeyBuilder.buildPopRetryTopic(requestHeader.getTopic(), requestHeader.getConsumerGroup()) : requestHeader.getTopic();
|
||||
long offset = getPopOffset(topic, requestHeader.getConsumerGroup(), queueId);
|
||||
|
||||
+18
@@ -33,6 +33,9 @@ public class NotificationRequestHeader extends TopicQueueRequestHeader {
|
||||
@CFNotNull
|
||||
private long bornTime;
|
||||
|
||||
private Boolean order = Boolean.FALSE;
|
||||
private String attemptId;
|
||||
|
||||
@CFNotNull
|
||||
@Override
|
||||
public void checkFields() throws RemotingCommandException {
|
||||
@@ -81,4 +84,19 @@ public class NotificationRequestHeader extends TopicQueueRequestHeader {
|
||||
this.queueId = queueId;
|
||||
}
|
||||
|
||||
public Boolean getOrder() {
|
||||
return order;
|
||||
}
|
||||
|
||||
public void setOrder(Boolean order) {
|
||||
this.order = order;
|
||||
}
|
||||
|
||||
public String getAttemptId() {
|
||||
return attemptId;
|
||||
}
|
||||
|
||||
public void setAttemptId(String attemptId) {
|
||||
this.attemptId = attemptId;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -172,12 +172,19 @@ public class RMQPopClient implements MQConsumer {
|
||||
|
||||
public CompletableFuture<Boolean> notification(String brokerAddr, String topic,
|
||||
String consumerGroup, int queueId, long pollTime, long bornTime, long timeoutMillis) {
|
||||
return notification(brokerAddr, topic, consumerGroup, queueId, null, null, pollTime, bornTime, timeoutMillis);
|
||||
}
|
||||
|
||||
public CompletableFuture<Boolean> notification(String brokerAddr, String topic,
|
||||
String consumerGroup, int queueId, Boolean order, String attemptId, long pollTime, long bornTime, long timeoutMillis) {
|
||||
NotificationRequestHeader requestHeader = new NotificationRequestHeader();
|
||||
requestHeader.setConsumerGroup(consumerGroup);
|
||||
requestHeader.setTopic(topic);
|
||||
requestHeader.setQueueId(queueId);
|
||||
requestHeader.setPollTime(pollTime);
|
||||
requestHeader.setBornTime(bornTime);
|
||||
requestHeader.setOrder(order);
|
||||
requestHeader.setAttemptId(attemptId);
|
||||
return this.mqClientAPI.notification(brokerAddr, requestHeader, timeoutMillis);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -66,6 +66,25 @@ public class NotificationIT extends BasePop {
|
||||
assertThat(result2).isFalse();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNotificationOrderly() throws Exception {
|
||||
long pollTime = 500;
|
||||
String attemptId = "attemptId";
|
||||
CompletableFuture<Boolean> future1 = client.notification(brokerAddr, topic, group, messageQueue.getQueueId(), true, attemptId, pollTime, System.currentTimeMillis(), 5000);
|
||||
CompletableFuture<Boolean> future2 = client.notification(brokerAddr, topic, group, messageQueue.getQueueId(), true, attemptId, pollTime, System.currentTimeMillis(), 5000);
|
||||
sendMessage(1);
|
||||
Boolean result1 = future1.get();
|
||||
assertThat(result1).isTrue();
|
||||
client.popMessageAsync(brokerAddr, messageQueue, 10000, 1, group, 1000, false,
|
||||
ConsumeInitMode.MIN, true, null, null, attemptId);
|
||||
Boolean result2 = future2.get();
|
||||
assertThat(result2).isTrue();
|
||||
|
||||
String attemptId2 = "attemptId2";
|
||||
CompletableFuture<Boolean> future3 = client.notification(brokerAddr, topic, group, messageQueue.getQueueId(), true, attemptId2, pollTime, System.currentTimeMillis(), 5000);
|
||||
assertThat(future3.get()).isFalse();
|
||||
}
|
||||
|
||||
protected void sendMessage(int num) {
|
||||
MessageQueueMsg mqMsgs = new MessageQueueMsg(Lists.newArrayList(messageQueue), num);
|
||||
producer.send(mqMsgs.getMsgsWithMQ());
|
||||
|
||||
Reference in New Issue
Block a user