From 2a119a75daa8e3fd349009c641325c48bf50be9d Mon Sep 17 00:00:00 2001 From: Hong-Yi Chen Date: Mon, 18 Aug 2025 21:36:08 +0800 Subject: [PATCH 1/4] KAFKA-19535: add integration tests for DescribeProducersOptions#brokerId --- .../DescribeProducersWithBrokerIdTest.java | 143 ++++++++++++++++++ 1 file changed, 143 insertions(+) create mode 100644 clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java diff --git a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java new file mode 100644 index 0000000000000..ccad82d908fad --- /dev/null +++ b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java @@ -0,0 +1,143 @@ +/* + * 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.kafka.clients.admin; + +import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.Node; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.NotLeaderOrFollowerException; +import org.apache.kafka.common.test.ClusterInstance; +import org.apache.kafka.common.test.api.ClusterConfigProperty; +import org.apache.kafka.common.test.api.ClusterTest; +import org.apache.kafka.common.test.api.ClusterTestDefaults; +import org.apache.kafka.test.TestUtils; + +import org.junit.jupiter.api.BeforeEach; + +import java.util.List; + +import static org.apache.kafka.coordinator.group.GroupCoordinatorConfig.GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG; +import static org.apache.kafka.coordinator.group.GroupCoordinatorConfig.OFFSETS_TOPIC_PARTITIONS_CONFIG; +import static org.apache.kafka.server.config.ServerLogConfigs.AUTO_CREATE_TOPICS_ENABLE_CONFIG; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; + + +@ClusterTestDefaults( + brokers = 4, + serverProperties = { + @ClusterConfigProperty(key = AUTO_CREATE_TOPICS_ENABLE_CONFIG, value = "false"), + @ClusterConfigProperty(key = OFFSETS_TOPIC_PARTITIONS_CONFIG, value = "1"), + @ClusterConfigProperty(key = GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG, value = "0") + } +) +class DescribeProducersWithBrokerIdTest { + private static final String TOPIC_NAME = "test-topic"; + private static final int NUM_PARTITIONS = 1; + private static final short REPLICATION_FACTOR = 3; + + private final ClusterInstance clusterInstance; + private final TopicPartition topicPartition; + + public DescribeProducersWithBrokerIdTest(ClusterInstance clusterInstance) { + this.clusterInstance = clusterInstance; + this.topicPartition = new TopicPartition(TOPIC_NAME, 0); + } + + private static void sendTestRecords(Producer producer) { + producer.send(new ProducerRecord<>(TOPIC_NAME, 0, "key-0".getBytes(), "value-0".getBytes())); + producer.flush(); + } + + @BeforeEach + void setUp() throws InterruptedException { + clusterInstance.createTopic(TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR); + } + + @ClusterTest + void testDescribeProducersDefaultRoutesToLeader() throws Exception { + try (Producer producer = clusterInstance.producer(); + var admin = clusterInstance.admin()) { + sendTestRecords(producer); + + var stateWithExplicitLeader = admin.describeProducers(List.of(topicPartition), + new DescribeProducersOptions().brokerId(clusterInstance.getLeaderBrokerId(topicPartition))) + .partitionResult(topicPartition).get(); + + var stateWithDefaultRouting = admin.describeProducers(List.of(topicPartition)) + .partitionResult(topicPartition).get(); + + assertNotNull(stateWithDefaultRouting); + assertFalse(stateWithDefaultRouting.activeProducers().isEmpty()); + assertEquals(stateWithExplicitLeader.activeProducers(), stateWithDefaultRouting.activeProducers()); + } + } + + @ClusterTest + void testDescribeProducersFromFollower() throws Exception { + try (Producer producer = clusterInstance.producer(); + var admin = clusterInstance.admin()) { + sendTestRecords(producer); + + var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); + var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() + .map(Node::id) + .toList(); + + var leaderBrokerId = clusterInstance.getLeaderBrokerId(topicPartition); + var followerBrokerId = replicaBrokerIds.stream() + .filter(id -> id != leaderBrokerId) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); + + var followerState = admin.describeProducers(List.of(topicPartition), + new DescribeProducersOptions().brokerId(followerBrokerId)) + .partitionResult(topicPartition).get(); + var leaderState = admin.describeProducers(List.of(topicPartition)) + .partitionResult(topicPartition).get(); + + assertNotNull(followerState); + assertFalse(followerState.activeProducers().isEmpty()); + assertEquals(leaderState.activeProducers(), followerState.activeProducers()); + } + } + + @ClusterTest + void testDescribeProducersWithInvalidBrokerId() throws Exception { + try (Producer producer = clusterInstance.producer(); + var admin = clusterInstance.admin()) { + sendTestRecords(producer); + + var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); + var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() + .map(Node::id) + .toList(); + + var nonReplicaBrokerId = clusterInstance.brokerIds().stream() + .filter(id -> !replicaBrokerIds.contains(id)) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); + + TestUtils.assertFutureThrows(NotLeaderOrFollowerException.class, + admin.describeProducers(List.of(topicPartition), + new DescribeProducersOptions().brokerId(nonReplicaBrokerId)) + .partitionResult(topicPartition)); + } + } +} \ No newline at end of file From 03708d3c7872299d9983eb4dc095dd4b4dd75eaf Mon Sep 17 00:00:00 2001 From: Hong-Yi Chen Date: Sat, 30 Aug 2025 13:20:10 +0800 Subject: [PATCH 2/4] Address comments --- .../DescribeProducersWithBrokerIdTest.java | 72 +++++++++---------- 1 file changed, 36 insertions(+), 36 deletions(-) diff --git a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java index ccad82d908fad..096629b4d1e4d 100644 --- a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java +++ b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java @@ -22,7 +22,6 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.NotLeaderOrFollowerException; import org.apache.kafka.common.test.ClusterInstance; -import org.apache.kafka.common.test.api.ClusterConfigProperty; import org.apache.kafka.common.test.api.ClusterTest; import org.apache.kafka.common.test.api.ClusterTestDefaults; import org.apache.kafka.test.TestUtils; @@ -31,21 +30,13 @@ import java.util.List; -import static org.apache.kafka.coordinator.group.GroupCoordinatorConfig.GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG; -import static org.apache.kafka.coordinator.group.GroupCoordinatorConfig.OFFSETS_TOPIC_PARTITIONS_CONFIG; -import static org.apache.kafka.server.config.ServerLogConfigs.AUTO_CREATE_TOPICS_ENABLE_CONFIG; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; @ClusterTestDefaults( - brokers = 4, - serverProperties = { - @ClusterConfigProperty(key = AUTO_CREATE_TOPICS_ENABLE_CONFIG, value = "false"), - @ClusterConfigProperty(key = OFFSETS_TOPIC_PARTITIONS_CONFIG, value = "1"), - @ClusterConfigProperty(key = GROUP_INITIAL_REBALANCE_DELAY_MS_CONFIG, value = "0") - } + brokers = 3 ) class DescribeProducersWithBrokerIdTest { private static final String TOPIC_NAME = "test-topic"; @@ -75,13 +66,15 @@ void testDescribeProducersDefaultRoutesToLeader() throws Exception { try (Producer producer = clusterInstance.producer(); var admin = clusterInstance.admin()) { sendTestRecords(producer); - - var stateWithExplicitLeader = admin.describeProducers(List.of(topicPartition), - new DescribeProducersOptions().brokerId(clusterInstance.getLeaderBrokerId(topicPartition))) - .partitionResult(topicPartition).get(); - var stateWithDefaultRouting = admin.describeProducers(List.of(topicPartition)) - .partitionResult(topicPartition).get(); + var stateWithExplicitLeader = admin.describeProducers( + List.of(topicPartition), + new DescribeProducersOptions().brokerId(clusterInstance.getLeaderBrokerId(topicPartition)) + ).partitionResult(topicPartition).get(); + + var stateWithDefaultRouting = admin.describeProducers( + List.of(topicPartition) + ).partitionResult(topicPartition).get(); assertNotNull(stateWithDefaultRouting); assertFalse(stateWithDefaultRouting.activeProducers().isEmpty()); @@ -97,20 +90,23 @@ void testDescribeProducersFromFollower() throws Exception { var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() - .map(Node::id) - .toList(); + .map(Node::id) + .toList(); var leaderBrokerId = clusterInstance.getLeaderBrokerId(topicPartition); var followerBrokerId = replicaBrokerIds.stream() - .filter(id -> id != leaderBrokerId) - .findFirst() - .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); + .filter(id -> id != leaderBrokerId) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); + + var followerState = admin.describeProducers( + List.of(topicPartition), + new DescribeProducersOptions().brokerId(followerBrokerId) + ).partitionResult(topicPartition).get(); - var followerState = admin.describeProducers(List.of(topicPartition), - new DescribeProducersOptions().brokerId(followerBrokerId)) - .partitionResult(topicPartition).get(); - var leaderState = admin.describeProducers(List.of(topicPartition)) - .partitionResult(topicPartition).get(); + var leaderState = admin.describeProducers( + List.of(topicPartition) + ).partitionResult(topicPartition).get(); assertNotNull(followerState); assertFalse(followerState.activeProducers().isEmpty()); @@ -118,26 +114,30 @@ void testDescribeProducersFromFollower() throws Exception { } } - @ClusterTest + @ClusterTest(brokers = 4) void testDescribeProducersWithInvalidBrokerId() throws Exception { try (Producer producer = clusterInstance.producer(); var admin = clusterInstance.admin()) { sendTestRecords(producer); - var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); + var topicDescription = admin.describeTopics( + List.of(topicPartition.topic()) + ).allTopicNames().get().get(topicPartition.topic()); + var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() - .map(Node::id) - .toList(); + .map(Node::id) + .toList(); var nonReplicaBrokerId = clusterInstance.brokerIds().stream() - .filter(id -> !replicaBrokerIds.contains(id)) - .findFirst() - .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); + .filter(id -> !replicaBrokerIds.contains(id)) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); TestUtils.assertFutureThrows(NotLeaderOrFollowerException.class, - admin.describeProducers(List.of(topicPartition), - new DescribeProducersOptions().brokerId(nonReplicaBrokerId)) - .partitionResult(topicPartition)); + admin.describeProducers( + List.of(topicPartition), + new DescribeProducersOptions().brokerId(nonReplicaBrokerId) + ).partitionResult(topicPartition)); } } } \ No newline at end of file From 107c9d5838c05e015a870180fea59194758ecae5 Mon Sep 17 00:00:00 2001 From: Hong-Yi Chen Date: Sat, 30 Aug 2025 15:06:08 +0800 Subject: [PATCH 3/4] refactor: extract helper method for brokerID retrieval --- .../DescribeProducersWithBrokerIdTest.java | 52 +++++++++---------- 1 file changed, 26 insertions(+), 26 deletions(-) diff --git a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java index 096629b4d1e4d..18d2c07b4e6d8 100644 --- a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java +++ b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java @@ -61,6 +61,30 @@ void setUp() throws InterruptedException { clusterInstance.createTopic(TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR); } + private List getReplicaBrokerIds(Admin admin, TopicPartition topicPartition) throws Exception { + var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); + return topicDescription.partitions().get(topicPartition.partition()).replicas().stream() + .map(Node::id) + .toList(); + } + + private int getNonReplicaBrokerId(Admin admin, TopicPartition topicPartition) throws Exception { + var replicaBrokerIds = getReplicaBrokerIds(admin, topicPartition); + return clusterInstance.brokerIds().stream() + .filter(id -> !replicaBrokerIds.contains(id)) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); + } + + private int getFollowerBrokerId(Admin admin, TopicPartition topicPartition) throws Exception { + var replicaBrokerIds = getReplicaBrokerIds(admin, topicPartition); + var leaderBrokerId = clusterInstance.getLeaderBrokerId(topicPartition); + return replicaBrokerIds.stream() + .filter(id -> id != leaderBrokerId) + .findFirst() + .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); + } + @ClusterTest void testDescribeProducersDefaultRoutesToLeader() throws Exception { try (Producer producer = clusterInstance.producer(); @@ -88,20 +112,9 @@ void testDescribeProducersFromFollower() throws Exception { var admin = clusterInstance.admin()) { sendTestRecords(producer); - var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); - var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() - .map(Node::id) - .toList(); - - var leaderBrokerId = clusterInstance.getLeaderBrokerId(topicPartition); - var followerBrokerId = replicaBrokerIds.stream() - .filter(id -> id != leaderBrokerId) - .findFirst() - .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); - var followerState = admin.describeProducers( List.of(topicPartition), - new DescribeProducersOptions().brokerId(followerBrokerId) + new DescribeProducersOptions().brokerId(getFollowerBrokerId(admin, topicPartition)) ).partitionResult(topicPartition).get(); var leaderState = admin.describeProducers( @@ -120,23 +133,10 @@ void testDescribeProducersWithInvalidBrokerId() throws Exception { var admin = clusterInstance.admin()) { sendTestRecords(producer); - var topicDescription = admin.describeTopics( - List.of(topicPartition.topic()) - ).allTopicNames().get().get(topicPartition.topic()); - - var replicaBrokerIds = topicDescription.partitions().get(topicPartition.partition()).replicas().stream() - .map(Node::id) - .toList(); - - var nonReplicaBrokerId = clusterInstance.brokerIds().stream() - .filter(id -> !replicaBrokerIds.contains(id)) - .findFirst() - .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); - TestUtils.assertFutureThrows(NotLeaderOrFollowerException.class, admin.describeProducers( List.of(topicPartition), - new DescribeProducersOptions().brokerId(nonReplicaBrokerId) + new DescribeProducersOptions().brokerId(getNonReplicaBrokerId(admin, topicPartition)) ).partitionResult(topicPartition)); } } From 5d7ccd463c1319f199374d0655d7b1436327a30f Mon Sep 17 00:00:00 2001 From: Hong-Yi Chen Date: Sun, 31 Aug 2025 13:37:49 +0800 Subject: [PATCH 4/4] Address comments --- .../DescribeProducersWithBrokerIdTest.java | 50 +++++++++---------- 1 file changed, 25 insertions(+), 25 deletions(-) diff --git a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java index 18d2c07b4e6d8..979989af1cafa 100644 --- a/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java +++ b/clients/clients-integration-tests/src/test/java/org/apache/kafka/clients/admin/DescribeProducersWithBrokerIdTest.java @@ -43,16 +43,16 @@ class DescribeProducersWithBrokerIdTest { private static final int NUM_PARTITIONS = 1; private static final short REPLICATION_FACTOR = 3; + private static final TopicPartition TOPIC_PARTITION = new TopicPartition(TOPIC_NAME, 0); + private final ClusterInstance clusterInstance; - private final TopicPartition topicPartition; public DescribeProducersWithBrokerIdTest(ClusterInstance clusterInstance) { this.clusterInstance = clusterInstance; - this.topicPartition = new TopicPartition(TOPIC_NAME, 0); } private static void sendTestRecords(Producer producer) { - producer.send(new ProducerRecord<>(TOPIC_NAME, 0, "key-0".getBytes(), "value-0".getBytes())); + producer.send(new ProducerRecord<>(TOPIC_NAME, TOPIC_PARTITION.partition(), "key-0".getBytes(), "value-0".getBytes())); producer.flush(); } @@ -61,28 +61,28 @@ void setUp() throws InterruptedException { clusterInstance.createTopic(TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR); } - private List getReplicaBrokerIds(Admin admin, TopicPartition topicPartition) throws Exception { - var topicDescription = admin.describeTopics(List.of(topicPartition.topic())).allTopicNames().get().get(topicPartition.topic()); - return topicDescription.partitions().get(topicPartition.partition()).replicas().stream() + private List getReplicaBrokerIds(Admin admin) throws Exception { + var topicDescription = admin.describeTopics(List.of(TOPIC_PARTITION.topic())).allTopicNames().get().get(TOPIC_PARTITION.topic()); + return topicDescription.partitions().get(TOPIC_PARTITION.partition()).replicas().stream() .map(Node::id) .toList(); } - private int getNonReplicaBrokerId(Admin admin, TopicPartition topicPartition) throws Exception { - var replicaBrokerIds = getReplicaBrokerIds(admin, topicPartition); + private int getNonReplicaBrokerId(Admin admin) throws Exception { + var replicaBrokerIds = getReplicaBrokerIds(admin); return clusterInstance.brokerIds().stream() .filter(id -> !replicaBrokerIds.contains(id)) .findFirst() .orElseThrow(() -> new IllegalStateException("No non-replica broker found")); } - private int getFollowerBrokerId(Admin admin, TopicPartition topicPartition) throws Exception { - var replicaBrokerIds = getReplicaBrokerIds(admin, topicPartition); - var leaderBrokerId = clusterInstance.getLeaderBrokerId(topicPartition); + private int getFollowerBrokerId(Admin admin) throws Exception { + var replicaBrokerIds = getReplicaBrokerIds(admin); + var leaderBrokerId = clusterInstance.getLeaderBrokerId(TOPIC_PARTITION); return replicaBrokerIds.stream() .filter(id -> id != leaderBrokerId) .findFirst() - .orElseThrow(() -> new IllegalStateException("No follower found for partition " + topicPartition)); + .orElseThrow(() -> new IllegalStateException("No follower found for partition " + TOPIC_PARTITION)); } @ClusterTest @@ -92,13 +92,13 @@ void testDescribeProducersDefaultRoutesToLeader() throws Exception { sendTestRecords(producer); var stateWithExplicitLeader = admin.describeProducers( - List.of(topicPartition), - new DescribeProducersOptions().brokerId(clusterInstance.getLeaderBrokerId(topicPartition)) - ).partitionResult(topicPartition).get(); + List.of(TOPIC_PARTITION), + new DescribeProducersOptions().brokerId(clusterInstance.getLeaderBrokerId(TOPIC_PARTITION)) + ).partitionResult(TOPIC_PARTITION).get(); var stateWithDefaultRouting = admin.describeProducers( - List.of(topicPartition) - ).partitionResult(topicPartition).get(); + List.of(TOPIC_PARTITION) + ).partitionResult(TOPIC_PARTITION).get(); assertNotNull(stateWithDefaultRouting); assertFalse(stateWithDefaultRouting.activeProducers().isEmpty()); @@ -113,13 +113,13 @@ void testDescribeProducersFromFollower() throws Exception { sendTestRecords(producer); var followerState = admin.describeProducers( - List.of(topicPartition), - new DescribeProducersOptions().brokerId(getFollowerBrokerId(admin, topicPartition)) - ).partitionResult(topicPartition).get(); + List.of(TOPIC_PARTITION), + new DescribeProducersOptions().brokerId(getFollowerBrokerId(admin)) + ).partitionResult(TOPIC_PARTITION).get(); var leaderState = admin.describeProducers( - List.of(topicPartition) - ).partitionResult(topicPartition).get(); + List.of(TOPIC_PARTITION) + ).partitionResult(TOPIC_PARTITION).get(); assertNotNull(followerState); assertFalse(followerState.activeProducers().isEmpty()); @@ -135,9 +135,9 @@ void testDescribeProducersWithInvalidBrokerId() throws Exception { TestUtils.assertFutureThrows(NotLeaderOrFollowerException.class, admin.describeProducers( - List.of(topicPartition), - new DescribeProducersOptions().brokerId(getNonReplicaBrokerId(admin, topicPartition)) - ).partitionResult(topicPartition)); + List.of(TOPIC_PARTITION), + new DescribeProducersOptions().brokerId(getNonReplicaBrokerId(admin)) + ).partitionResult(TOPIC_PARTITION)); } } } \ No newline at end of file