GH-2525: Fix Paused Partition Resume (Rebalance)
Resolves https://github.com/spring-projects/spring-kafka/issues/2525 If a paused partition is revoked and re-assigned during a retry topic pause delay, the partition is not resumed because it no longer exists in the `pausedPartitions` field. When re-pausing a paused partition, it must be re-added to the paused partitions so that it will not be filtered out of the resume logic. **cherry-pick to 2.9.x**
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2016-2022 the original author or authors.
|
||||
* Copyright 2016-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -3531,6 +3531,7 @@ public class KafkaMessageListenerContainer<K, V> // NOSONAR line count
|
||||
ListenerConsumer.this.consumer.pause(toRepause);
|
||||
ListenerConsumer.this.logger.debug(() -> "Paused consumption from: " + toRepause);
|
||||
publishConsumerPausedEvent(toRepause, "Re-paused after rebalance");
|
||||
ListenerConsumer.this.pausedPartitions.addAll(toRepause);
|
||||
}
|
||||
this.revoked.removeAll(toRepause);
|
||||
ListenerConsumer.this.pausedPartitions.removeAll(this.revoked);
|
||||
|
||||
@@ -2842,6 +2842,95 @@ public class KafkaMessageListenerContainerTests {
|
||||
container.stop();
|
||||
}
|
||||
|
||||
@SuppressWarnings({ "unchecked" })
|
||||
@Test
|
||||
public void resumePartitionAfterRevokeAndReAssign() throws Exception {
|
||||
ConsumerFactory<Integer, String> cf = mock(ConsumerFactory.class);
|
||||
Consumer<Integer, String> consumer = mock(Consumer.class);
|
||||
given(cf.createConsumer(eq("grp"), eq("clientId"), isNull(), any())).willReturn(consumer);
|
||||
AtomicBoolean first = new AtomicBoolean(true);
|
||||
TopicPartition tp0 = new TopicPartition("foo", 0);
|
||||
TopicPartition tp1 = new TopicPartition("foo", 1);
|
||||
given(consumer.assignment()).willReturn(Set.of(tp0, tp1));
|
||||
final CountDownLatch pauseLatch1 = new CountDownLatch(1);
|
||||
final CountDownLatch suspendConsumerThread = new CountDownLatch(1);
|
||||
Set<TopicPartition> pausedParts = ConcurrentHashMap.newKeySet();
|
||||
Thread testThread = Thread.currentThread();
|
||||
AtomicBoolean paused = new AtomicBoolean();
|
||||
willAnswer(i -> {
|
||||
pausedParts.clear();
|
||||
pausedParts.addAll(i.getArgument(0));
|
||||
if (!Thread.currentThread().equals(testThread)) {
|
||||
paused.set(true);
|
||||
}
|
||||
return null;
|
||||
}).given(consumer).pause(any());
|
||||
given(consumer.paused()).willReturn(pausedParts);
|
||||
given(consumer.poll(any(Duration.class))).willAnswer(i -> {
|
||||
if (paused.get()) {
|
||||
pauseLatch1.countDown();
|
||||
// hold up the consumer thread while we revoke/assign partitions on the test thread
|
||||
suspendConsumerThread.await(10, TimeUnit.SECONDS);
|
||||
}
|
||||
Thread.sleep(50);
|
||||
return ConsumerRecords.empty();
|
||||
});
|
||||
AtomicReference<ConsumerRebalanceListener> rebal = new AtomicReference<>();
|
||||
Collection<String> foos = new ArrayList<>();
|
||||
foos.add("foo");
|
||||
willAnswer(inv -> {
|
||||
rebal.set(inv.getArgument(1));
|
||||
rebal.get().onPartitionsAssigned(Set.of(tp0, tp1));
|
||||
return null;
|
||||
}).given(consumer).subscribe(eq(foos), any(ConsumerRebalanceListener.class));
|
||||
final CountDownLatch resumeLatch = new CountDownLatch(1);
|
||||
willAnswer(i -> {
|
||||
pausedParts.removeAll(i.getArgument(0));
|
||||
resumeLatch.countDown();
|
||||
return null;
|
||||
}).given(consumer).resume(any());
|
||||
ContainerProperties containerProps = new ContainerProperties("foo");
|
||||
containerProps.setGroupId("grp");
|
||||
containerProps.setAckMode(AckMode.RECORD);
|
||||
containerProps.setClientId("clientId");
|
||||
containerProps.setIdleEventInterval(100L);
|
||||
containerProps.setMessageListener((MessageListener) rec -> { });
|
||||
containerProps.setMissingTopicsFatal(false);
|
||||
KafkaMessageListenerContainer<Integer, String> container =
|
||||
new KafkaMessageListenerContainer<>(cf, containerProps);
|
||||
container.start();
|
||||
container.pausePartition(tp0);
|
||||
container.pausePartition(tp1);
|
||||
assertThat(pauseLatch1.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
assertThat(pausedParts).hasSize(2)
|
||||
.contains(tp0, tp1);
|
||||
rebal.get().onPartitionsRevoked(Set.of(tp0, tp1));
|
||||
rebal.get().onPartitionsAssigned(Collections.singleton(tp0));
|
||||
rebal.get().onPartitionsRevoked(Set.of(tp0));
|
||||
rebal.get().onPartitionsAssigned(Set.of(tp0, tp1));
|
||||
assertThat(pausedParts).hasSize(2)
|
||||
.contains(tp0, tp1);
|
||||
assertThat(container).extracting("listenerConsumer")
|
||||
.extracting("pausedPartitions")
|
||||
.asInstanceOf(InstanceOfAssertFactories.collection(TopicPartition.class))
|
||||
.hasSize(2)
|
||||
.contains(tp0, tp1);
|
||||
assertThat(container)
|
||||
.extracting("pauseRequestedPartitions")
|
||||
.asInstanceOf(InstanceOfAssertFactories.collection(TopicPartition.class))
|
||||
.hasSize(2)
|
||||
.contains(tp0, tp1);
|
||||
container.resumePartition(tp0);
|
||||
container.resumePartition(tp1);
|
||||
suspendConsumerThread.countDown();
|
||||
assertThat(resumeLatch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
assertThat(pausedParts).hasSize(0);
|
||||
ArgumentCaptor<List<TopicPartition>> resumed = ArgumentCaptor.forClass(List.class);
|
||||
verify(consumer).resume(resumed.capture());
|
||||
assertThat(resumed.getValue()).contains(tp0, tp1);
|
||||
container.stop();
|
||||
}
|
||||
|
||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||
@Test
|
||||
public void testInitialSeek() throws Exception {
|
||||
|
||||
Reference in New Issue
Block a user