Add and refactor integration tests in DefaultShareConsumerFactoryTests

* Add and refactor integration tests in DefaultShareConsumerFactoryTests

- Added tests to check that multiple shared consumers each get records and all data is consumed.
- Refactored shared test logic into a helper method.
- Other cleanup in the test

Signed-off-by: Soby Chacko <soby.chacko@broadcom.com>
This commit is contained in:
Soby Chacko
2025-06-03 18:12:53 -04:00
committed by GitHub
parent bcb5395fa2
commit 9674334eda

View File

@@ -17,16 +17,21 @@
package org.springframework.kafka.core;
import java.time.Duration;
import java.util.Arrays;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AlterConfigOp;
import org.apache.kafka.clients.admin.ConfigEntry;
import org.apache.kafka.clients.consumer.AcknowledgeType;
import org.apache.kafka.clients.consumer.ShareConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
@@ -49,7 +54,7 @@ import static org.assertj.core.api.Assertions.assertThatThrownBy;
* @since 4.0
*/
@EmbeddedKafka(
topics = {"embedded-share-test"}, partitions = 1,
topics = {"embedded-share-test", "embedded-share-distribution-test"}, partitions = 1,
brokerProperties = {
"unstable.api.versions.enable=true",
"group.coordinator.rebalance.protocols=classic,share",
@@ -144,7 +149,6 @@ class DefaultShareConsumerFactoryTests {
}
@Test
@SuppressWarnings("try")
void integrationTestDefaultShareConsumerFactory(EmbeddedKafkaBroker broker) throws Exception {
final String topic = "embedded-share-test";
final String groupId = "testGroup";
@@ -159,23 +163,7 @@ class DefaultShareConsumerFactoryTests {
producer.send(new ProducerRecord<>(topic, "key", "integration-test-value")).get();
}
Map<String, Object> adminProperties = new HashMap<>();
adminProperties.put("bootstrap.servers", bootstrapServers);
// For this test: force new share groups to start from the beginning of the topic.
// This is NOT the same as the usual consumer auto.offset.reset; it's a group config,
// so use AdminClient to set share.auto.offset.reset = earliest for our test group.
try (AdminClient ignored = AdminClient.create(adminProperties)) {
ConfigEntry entry = new ConfigEntry("share.auto.offset.reset", "earliest");
AlterConfigOp op = new AlterConfigOp(entry, AlterConfigOp.OpType.SET);
Map<ConfigResource, Collection<AlterConfigOp>> configs = Map.of(
new ConfigResource(ConfigResource.Type.GROUP, "testGroup"), Arrays.asList(op));
try (Admin admin = AdminClient.create(adminProperties)) {
admin.incrementalAlterConfigs(configs).all().get();
}
}
setShareAutoOffsetResetEarliest(bootstrapServers, groupId);
var consumerProps = new HashMap<String, Object>();
consumerProps.put("bootstrap.servers", bootstrapServers);
@@ -197,4 +185,101 @@ class DefaultShareConsumerFactoryTests {
consumer.close();
}
@Test
void integrationTestSharedConsumersDistribution(EmbeddedKafkaBroker broker) throws Exception {
String topic = "shared-consumer-dist-test";
final String groupId = "distributionTestGroup";
int recordCount = 8;
List<String> consumerIds = List.of("client-dist-1", "client-dist-2");
List<String> allReceived = runSharedConsumerTest(topic, groupId, consumerIds, recordCount, broker);
// Assert all records were received (no loss and no duplicates)
assertThat(allReceived)
.containsExactlyInAnyOrder(
topic + "-value-0",
topic + "-value-1",
topic + "-value-2",
topic + "-value-3",
topic + "-value-4",
topic + "-value-5",
topic + "-value-6",
topic + "-value-7"
)
.doesNotHaveDuplicates();
}
/**
* Runs multiple Kafka consumers in parallel using ExecutorService, collects all records received,
* and returns a list of all record values received by all consumers.
*/
private static List<String> runSharedConsumerTest(String topic, String groupId,
List<String> consumerIds, int recordCount, EmbeddedKafkaBroker broker) throws Exception {
var bootstrapServers = broker.getBrokersAsString();
var producerProps = new java.util.Properties();
producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
try (var producer = new KafkaProducer<String, String>(producerProps)) {
for (int i = 0; i < recordCount; i++) {
producer.send(new ProducerRecord<>(topic, "key" + i, topic + "-value-" + i)).get();
}
producer.flush();
}
setShareAutoOffsetResetEarliest(bootstrapServers, groupId);
List<String> allReceived = Collections.synchronizedList(new ArrayList<>());
var latch = new java.util.concurrent.CountDownLatch(recordCount);
ExecutorService executor = Executors.newCachedThreadPool();
DefaultShareConsumerFactory<String, String> shareConsumerFactory = new DefaultShareConsumerFactory<>(
Map.of(
"bootstrap.servers", bootstrapServers,
"key.deserializer", org.apache.kafka.common.serialization.StringDeserializer.class,
"value.deserializer", org.apache.kafka.common.serialization.StringDeserializer.class
)
);
for (int i = 0; i < consumerIds.size(); i++) {
final int idx = i;
executor.submit(() -> {
try (var consumer = shareConsumerFactory.createShareConsumer(groupId, consumerIds.get(idx))) {
consumer.subscribe(Collections.singletonList(topic));
while (latch.getCount() > 0) {
var records = consumer.poll(Duration.ofMillis(200));
for (var r : records) {
allReceived.add(r.value());
consumer.acknowledge(r, AcknowledgeType.ACCEPT);
latch.countDown();
}
}
}
});
}
assertThat(latch.await(10, TimeUnit.SECONDS))
.as("All records should be received within timeout")
.isTrue();
executor.shutdown();
assertThat(executor.awaitTermination(10, TimeUnit.SECONDS))
.as("Executor should terminate after shutdown")
.isTrue();
return allReceived;
}
/**
* Sets the share.auto.offset.reset group config to earliest for the given groupId,
* using the provided bootstrapServers.
*/
private static void setShareAutoOffsetResetEarliest(String bootstrapServers, String groupId) throws Exception {
Map<String, Object> adminProperties = new HashMap<>();
adminProperties.put("bootstrap.servers", bootstrapServers);
ConfigEntry entry = new ConfigEntry("share.auto.offset.reset", "earliest");
AlterConfigOp op = new AlterConfigOp(entry, AlterConfigOp.OpType.SET);
Map<ConfigResource, Collection<AlterConfigOp>> configs = Map.of(
new ConfigResource(ConfigResource.Type.GROUP, groupId), List.of(op));
try (Admin admin = AdminClient.create(adminProperties)) {
admin.incrementalAlterConfigs(configs).all().get();
}
}
}