mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] Refactor partition generation and add unit test.
This commit is contained in:
@@ -17,47 +17,9 @@
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.service;
|
||||
|
||||
import apache.rocketmq.v1.AckMessageRequest;
|
||||
import apache.rocketmq.v1.AckMessageResponse;
|
||||
import apache.rocketmq.v1.ChangeInvisibleDurationRequest;
|
||||
import apache.rocketmq.v1.ChangeInvisibleDurationResponse;
|
||||
import apache.rocketmq.v1.EndTransactionRequest;
|
||||
import apache.rocketmq.v1.EndTransactionResponse;
|
||||
import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest;
|
||||
import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse;
|
||||
import apache.rocketmq.v1.HealthCheckRequest;
|
||||
import apache.rocketmq.v1.HealthCheckResponse;
|
||||
import apache.rocketmq.v1.HeartbeatRequest;
|
||||
import apache.rocketmq.v1.HeartbeatResponse;
|
||||
import apache.rocketmq.v1.NackMessageRequest;
|
||||
import apache.rocketmq.v1.NackMessageResponse;
|
||||
import apache.rocketmq.v1.NoopCommand;
|
||||
import apache.rocketmq.v1.NotifyClientTerminationRequest;
|
||||
import apache.rocketmq.v1.NotifyClientTerminationResponse;
|
||||
import apache.rocketmq.v1.PollCommandRequest;
|
||||
import apache.rocketmq.v1.PollCommandResponse;
|
||||
import apache.rocketmq.v1.PullMessageRequest;
|
||||
import apache.rocketmq.v1.PullMessageResponse;
|
||||
import apache.rocketmq.v1.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v1.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v1.QueryOffsetRequest;
|
||||
import apache.rocketmq.v1.QueryOffsetResponse;
|
||||
import apache.rocketmq.v1.QueryRouteRequest;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.ReceiveMessageRequest;
|
||||
import apache.rocketmq.v1.ReceiveMessageResponse;
|
||||
import apache.rocketmq.v1.ReportMessageConsumptionResultRequest;
|
||||
import apache.rocketmq.v1.ReportMessageConsumptionResultResponse;
|
||||
import apache.rocketmq.v1.ReportThreadStackTraceRequest;
|
||||
import apache.rocketmq.v1.ReportThreadStackTraceResponse;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import apache.rocketmq.v1.SendMessageRequest;
|
||||
import apache.rocketmq.v1.SendMessageResponse;
|
||||
import apache.rocketmq.v1.*;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
@@ -69,15 +31,14 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ReceiveMessageService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService;
|
||||
import org.apache.rocketmq.proxy.grpc.service.cluster.*;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
|
||||
public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService {
|
||||
private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
|
||||
@@ -85,7 +46,6 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc
|
||||
new ThreadFactoryImpl("ClusterGrpcServiceScheduledThread"));
|
||||
|
||||
private final ChannelManager channelManager;
|
||||
|
||||
private final ConnectorManager connectorManager;
|
||||
private final ProducerService producerService;
|
||||
private final ReceiveMessageService receiveMessageService;
|
||||
|
||||
+31
-27
@@ -16,33 +16,25 @@
|
||||
*/
|
||||
package org.apache.rocketmq.proxy.grpc.service.cluster;
|
||||
|
||||
import apache.rocketmq.v1.Assignment;
|
||||
import apache.rocketmq.v1.Broker;
|
||||
import apache.rocketmq.v1.Endpoints;
|
||||
import apache.rocketmq.v1.Partition;
|
||||
import apache.rocketmq.v1.Permission;
|
||||
import apache.rocketmq.v1.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v1.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v1.QueryRouteRequest;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import apache.rocketmq.v1.*;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.constant.PermName;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.connector.ConnectorManager;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
|
||||
import org.apache.rocketmq.proxy.connector.route.TopicRouteHelper;
|
||||
import org.apache.rocketmq.proxy.grpc.common.Converter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ParameterConverter;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseHook;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
|
||||
public class RouteService extends BaseService {
|
||||
|
||||
private volatile ParameterConverter<Endpoints, Endpoints> queryRouteEndpointConverter;
|
||||
@@ -148,21 +140,33 @@ public class RouteService extends BaseService {
|
||||
r = queueData.getReadQueueNums();
|
||||
}
|
||||
|
||||
for (int i = 0; i < (rw + r + w); i++) {
|
||||
Partition.Builder builder = Partition.newBuilder()
|
||||
// r here means readOnly queue nums, w means writeOnly queue nums, while rw means readable and writable queue nums.
|
||||
int queueIdIndex = 0;
|
||||
for(int i = 0; i < r; i++){
|
||||
Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ);
|
||||
partitionList.add(partition);
|
||||
}
|
||||
|
||||
for(int i = 0; i < w; i++){
|
||||
Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.WRITE);
|
||||
partitionList.add(partition);
|
||||
}
|
||||
|
||||
for (int i = 0; i < rw; i++) {
|
||||
Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ_WRITE);
|
||||
partitionList.add(partition);
|
||||
}
|
||||
|
||||
return partitionList;
|
||||
}
|
||||
|
||||
private static Partition buildPartition(Broker broker, Resource topic, int queueId, Permission perm) {
|
||||
Partition.Builder builder = Partition.newBuilder()
|
||||
.setBroker(broker)
|
||||
.setTopic(topic)
|
||||
.setId(i);
|
||||
if (i < r) {
|
||||
builder.setPermission(Permission.READ);
|
||||
} else if (i < w) {
|
||||
builder.setPermission(Permission.WRITE);
|
||||
} else {
|
||||
builder.setPermission(Permission.READ_WRITE);
|
||||
}
|
||||
partitionList.add(builder.build());
|
||||
}
|
||||
return partitionList;
|
||||
.setId(queueId);
|
||||
builder.setPermission(perm);
|
||||
return builder.build();
|
||||
}
|
||||
|
||||
public CompletableFuture<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request) {
|
||||
|
||||
+93
-164
@@ -1,177 +1,106 @@
|
||||
/*
|
||||
* Licensed to the Apache Software Foundation (ASF) under one or more
|
||||
* contributor license agreements. See the NOTICE file distributed with
|
||||
* this work for additional information regarding copyright ownership.
|
||||
* The ASF licenses this file to You under the Apache License, Version 2.0
|
||||
* (the "License"); you may not use this file except in compliance with
|
||||
* the License. You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.apache.rocketmq.proxy.grpc.service.cluster;
|
||||
|
||||
import apache.rocketmq.v1.Address;
|
||||
import apache.rocketmq.v1.AddressScheme;
|
||||
import apache.rocketmq.v1.Endpoints;
|
||||
import apache.rocketmq.v1.QueryAssignmentRequest;
|
||||
import apache.rocketmq.v1.QueryAssignmentResponse;
|
||||
import apache.rocketmq.v1.QueryRouteRequest;
|
||||
import apache.rocketmq.v1.QueryRouteResponse;
|
||||
import apache.rocketmq.v1.Broker;
|
||||
import apache.rocketmq.v1.Partition;
|
||||
import apache.rocketmq.v1.Permission;
|
||||
import apache.rocketmq.v1.Resource;
|
||||
import com.google.rpc.Code;
|
||||
import io.grpc.Context;
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.client.exception.MQClientException;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.route.BrokerData;
|
||||
import org.apache.rocketmq.common.constant.PermName;
|
||||
import org.apache.rocketmq.common.protocol.route.QueueData;
|
||||
import org.apache.rocketmq.common.protocol.route.TopicRouteData;
|
||||
import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper;
|
||||
import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.mockito.Mockito.when;
|
||||
import java.util.List;
|
||||
|
||||
public class RouteServiceTest extends BaseServiceTest {
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
@Override
|
||||
public void beforeEach() throws Throwable {
|
||||
TopicRouteData routeData = new TopicRouteData();
|
||||
public class RouteServiceTest {
|
||||
public static final String BROKER_NAME = "brokerName";
|
||||
public static final String NAMESPACE = "namespace";
|
||||
public static final String TOPIC = "topic";
|
||||
public static final Broker MOCK_BROKER = Broker.newBuilder().setName(BROKER_NAME).build();
|
||||
public static final Resource MOCK_TOPIC = Resource.newBuilder()
|
||||
.setName(TOPIC)
|
||||
.setResourceNamespace(NAMESPACE)
|
||||
.build();
|
||||
|
||||
List<BrokerData> brokerDataList = new ArrayList<>();
|
||||
BrokerData brokerData = new BrokerData();
|
||||
brokerData.setCluster("cluster");
|
||||
brokerData.setBrokerName("brokerName");
|
||||
HashMap<Long, String> brokerAddrs = new HashMap<Long, String>() {{
|
||||
put(0L, "127.0.0.1:10911");
|
||||
}};
|
||||
brokerData.setBrokerAddrs(brokerAddrs);
|
||||
brokerDataList.add(brokerData);
|
||||
@Before
|
||||
public void before() throws Exception {
|
||||
}
|
||||
|
||||
List<QueueData> queueDataList = new ArrayList<>();
|
||||
@After
|
||||
public void after() throws Exception {
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testGenPartitionFromQueueData() throws Exception {
|
||||
// test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues.
|
||||
QueueData queueDataWith8R8WPermRW = mockQueueData(8, 8, PermName.PERM_READ | PermName.PERM_WRITE);
|
||||
List<Partition> partitionWith8R8WPermRW = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermRW, MOCK_TOPIC, MOCK_BROKER);
|
||||
assertThat(partitionWith8R8WPermRW.size()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0);
|
||||
assertThat(partitionWith8R8WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0);
|
||||
|
||||
// test queueData with 8 read queues, 8 write queues, and read only permission, expect 8 read only queues.
|
||||
QueueData queueDataWith8R8WPermR = mockQueueData(8, 8, PermName.PERM_READ);
|
||||
List<Partition> partitionWith8R8WPermR = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermR, MOCK_TOPIC, MOCK_BROKER);
|
||||
assertThat(partitionWith8R8WPermR.size()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0);
|
||||
assertThat(partitionWith8R8WPermR.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0);
|
||||
|
||||
// test queueData with 8 read queues, 8 write queues, and write only permission, expect 8 write only queues.
|
||||
QueueData queueDataWith8R8WPermW = mockQueueData(8, 8, PermName.PERM_WRITE);
|
||||
List<Partition> partitionWith8R8WPermW = RouteService.genPartitionFromQueueData(queueDataWith8R8WPermW, MOCK_TOPIC, MOCK_BROKER);
|
||||
assertThat(partitionWith8R8WPermW.size()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(8);
|
||||
assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0);
|
||||
assertThat(partitionWith8R8WPermW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0);
|
||||
|
||||
// test queueData with 8 read queues, 0 write queues, and rw permission, expect 8 read only queues.
|
||||
QueueData queueDataWith8R0WPermRW = mockQueueData(8, 0, PermName.PERM_READ | PermName.PERM_WRITE);
|
||||
List<Partition> partitionWith8R0WPermRW = RouteService.genPartitionFromQueueData(queueDataWith8R0WPermRW, MOCK_TOPIC, MOCK_BROKER);
|
||||
assertThat(partitionWith8R0WPermRW.size()).isEqualTo(8);
|
||||
assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(8);
|
||||
assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(0);
|
||||
assertThat(partitionWith8R0WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(0);
|
||||
|
||||
// test queueData with 4 read queues, 8 write queues, and rw permission, expect 4 rw queues and 4 write only queues.
|
||||
QueueData queueDataWith4R8WPermRW = mockQueueData(4, 8, PermName.PERM_READ | PermName.PERM_WRITE);
|
||||
List<Partition> partitionWith4R8WPermRW = RouteService.genPartitionFromQueueData(queueDataWith4R8WPermRW, MOCK_TOPIC, MOCK_BROKER);
|
||||
assertThat(partitionWith4R8WPermRW.size()).isEqualTo(8);
|
||||
assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.WRITE).count()).isEqualTo(4);
|
||||
assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ_WRITE).count()).isEqualTo(4);
|
||||
assertThat(partitionWith4R8WPermRW.stream().filter(a -> a.getPermission() == Permission.READ).count()).isEqualTo(0);
|
||||
|
||||
}
|
||||
|
||||
private QueueData mockQueueData(int r, int w, int perm) {
|
||||
QueueData queueData = new QueueData();
|
||||
queueData.setPerm(6);
|
||||
queueData.setWriteQueueNums(8);
|
||||
queueData.setReadQueueNums(8);
|
||||
queueData.setBrokerName("brokerName");
|
||||
queueDataList.add(queueData);
|
||||
|
||||
routeData.setBrokerDatas(brokerDataList);
|
||||
routeData.setQueueDatas(queueDataList);
|
||||
|
||||
MessageQueueWrapper messageQueueWrapper = new MessageQueueWrapper("topic", routeData);
|
||||
when(this.topicRouteCache.getMessageQueue("topic")).thenReturn(messageQueueWrapper);
|
||||
|
||||
when(this.topicRouteCache.getMessageQueue("notExistTopic")).thenThrow(new MQClientException(ResponseCode.TOPIC_NOT_EXIST, ""));
|
||||
queueData.setBrokerName(BROKER_NAME);
|
||||
queueData.setReadQueueNums(r);
|
||||
queueData.setWriteQueueNums(w);
|
||||
queueData.setPerm(perm);
|
||||
return queueData;
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWithInvalidEndpoints() {
|
||||
RouteService routeService = new RouteService(this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRoute() {
|
||||
RouteService routeService = new RouteService(this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(8, response.getPartitionsCount());
|
||||
assertEquals("host", response.getPartitions(0).getBroker()
|
||||
.getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryRouteWhenTopicNotExist() {
|
||||
RouteService routeService = new RouteService(this.clientManager);
|
||||
|
||||
CompletableFuture<QueryRouteResponse> future = routeService.queryRoute(Context.current(), QueryRouteRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("notExistTopic")
|
||||
.build())
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryRouteResponse response = future.get();
|
||||
assertEquals(Code.NOT_FOUND.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryAssignmentInvalidEndpoints() {
|
||||
RouteService routeService = new RouteService(this.clientManager);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.INVALID_ARGUMENT.getNumber(), response.getCommon().getStatus().getCode());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testQueryAssignment() {
|
||||
RouteService routeService = new RouteService(this.clientManager);
|
||||
|
||||
CompletableFuture<QueryAssignmentResponse> future = routeService.queryAssignment(Context.current(), QueryAssignmentRequest.newBuilder()
|
||||
.setEndpoints(Endpoints.newBuilder()
|
||||
.addAddresses(Address.newBuilder()
|
||||
.setPort(80)
|
||||
.setHost("host")
|
||||
.build())
|
||||
.setScheme(AddressScheme.DOMAIN_NAME)
|
||||
.build())
|
||||
.setTopic(Resource.newBuilder()
|
||||
.setName("topic")
|
||||
.build())
|
||||
.setGroup(Resource.newBuilder()
|
||||
.setName("group")
|
||||
.build())
|
||||
.setClientId("clientId")
|
||||
.build());
|
||||
|
||||
try {
|
||||
QueryAssignmentResponse response = future.get();
|
||||
assertEquals(Code.OK.getNumber(), response.getCommon().getStatus().getCode());
|
||||
assertEquals(1, response.getAssignmentsCount());
|
||||
assertEquals("brokerName", response.getAssignments(0).getPartition().getBroker().getName());
|
||||
assertEquals("host", response.getAssignments(0).getPartition().getBroker().getEndpoints().getAddresses(0).getHost());
|
||||
} catch (Exception e) {
|
||||
assertNull(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user