mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-01 05:46:12 +08:00
[ISSUE #9928] Add Priority IT for GRPC protocol
This commit is contained in:
@@ -118,4 +118,9 @@ public class ClusterGrpcIT extends GrpcBaseIT {
|
||||
public void testConsumeOrderly() throws Exception {
|
||||
super.testConsumeOrderly();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimpleConsumerSendAndRecvPriorityMessage() throws Exception {
|
||||
super.testSimpleConsumerSendAndRecvPriorityMessage();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -629,6 +629,56 @@ public class GrpcBaseIT extends BaseConf {
|
||||
}
|
||||
}
|
||||
|
||||
public void testSimpleConsumerSendAndRecvPriorityMessage() throws Exception {
|
||||
String topic = initTopicOnSampleTopicBroker(BROKER1_NAME, TopicMessageType.PRIORITY);
|
||||
String group = MQRandomUtils.getRandomConsumerGroup();
|
||||
|
||||
// init consumer offset
|
||||
this.sendClientSettings(stub, buildSimpleConsumerClientSettings(group)).get();
|
||||
receiveMessage(blockingStub, topic, group, 1);
|
||||
|
||||
this.sendClientSettings(stub, buildProducerClientSettings(topic)).get();
|
||||
for (int i = 0; i < BaseConf.QUEUE_NUMBERS; i++) {
|
||||
String messageId = createUniqID();
|
||||
SendMessageResponse sendResponse = blockingStub.sendMessage(SendMessageRequest.newBuilder()
|
||||
.addMessages(Message.newBuilder()
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName(topic)
|
||||
.build())
|
||||
.setSystemProperties(SystemProperties.newBuilder()
|
||||
.setMessageId(messageId)
|
||||
.setQueueId(0)
|
||||
.setMessageType(MessageType.PRIORITY)
|
||||
.setBodyEncoding(Encoding.GZIP)
|
||||
.setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis()))
|
||||
.setBornHost(StringUtils.defaultString(NetworkUtil.getLocalAddress(), "127.0.0.1:1234"))
|
||||
.setPriority(i)
|
||||
.build())
|
||||
.setBody(ByteString.copyFromUtf8("hello"))
|
||||
.build())
|
||||
.build());
|
||||
assertSendMessage(sendResponse, messageId);
|
||||
}
|
||||
|
||||
this.sendClientSettings(stub, buildSimpleConsumerClientSettings(group)).get();
|
||||
List<Message> recvList = new ArrayList<>();
|
||||
try {
|
||||
await().atMost(java.time.Duration.ofSeconds(10)).until(() -> {
|
||||
List<Message> messageList = getMessageFromReceiveMessageResponse(receiveMessage(blockingStub, topic, group));
|
||||
if (messageList.isEmpty()) {
|
||||
return false;
|
||||
}
|
||||
recvList.addAll(messageList);
|
||||
return recvList.size() == BaseConf.QUEUE_NUMBERS;
|
||||
});
|
||||
} catch (Exception e) {
|
||||
}
|
||||
for (int i = 0; i < BaseConf.QUEUE_NUMBERS; i++) {
|
||||
// default priority order: 0 as lowest priority
|
||||
assertThat(recvList.get(i).getSystemProperties().getPriority()).isEqualTo(BaseConf.QUEUE_NUMBERS - i - 1);
|
||||
}
|
||||
}
|
||||
|
||||
public List<ReceiveMessageResponse> receiveMessage(MessagingServiceGrpc.MessagingServiceBlockingStub stub,
|
||||
String topic, String group) {
|
||||
return receiveMessage(stub, topic, group, 15);
|
||||
|
||||
@@ -106,4 +106,9 @@ public class LocalGrpcIT extends GrpcBaseIT {
|
||||
public void testConsumeOrderly() throws Exception {
|
||||
super.testConsumeOrderly();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSimpleConsumerSendAndRecvPriorityMessage() throws Exception {
|
||||
super.testSimpleConsumerSendAndRecvPriorityMessage();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user