Upgrade to spring-kafka 2.6.0
This commit is contained in:
@@ -98,7 +98,7 @@ ext {
|
|||||||
soapVersion = '1.4.0'
|
soapVersion = '1.4.0'
|
||||||
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.3.0-SNAPSHOT'
|
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.3.0-SNAPSHOT'
|
||||||
springDataVersion = project.hasProperty('springDataVersion') ? project.springDataVersion : '2020.0.0-SNAPSHOT'
|
springDataVersion = project.hasProperty('springDataVersion') ? project.springDataVersion : '2020.0.0-SNAPSHOT'
|
||||||
springKafkaVersion = '2.5.4.BUILD-SNAPSHOT'
|
springKafkaVersion = '2.6.0-SNAPSHOT'
|
||||||
springSecurityVersion = project.hasProperty('springSecurityVersion') ? project.springSecurityVersion : '5.4.0-SNAPSHOT'
|
springSecurityVersion = project.hasProperty('springSecurityVersion') ? project.springSecurityVersion : '5.4.0-SNAPSHOT'
|
||||||
springRetryVersion = '1.3.0'
|
springRetryVersion = '1.3.0'
|
||||||
springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.3.0-SNAPSHOT'
|
springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.3.0-SNAPSHOT'
|
||||||
|
|||||||
@@ -101,6 +101,8 @@ public class KafkaMessageSource<K, V> extends AbstractMessageSource<Object> impl
|
|||||||
|
|
||||||
private static final long MIN_ASSIGN_TIMEOUT = 2000L;
|
private static final long MIN_ASSIGN_TIMEOUT = 2000L;
|
||||||
|
|
||||||
|
private static int DEFAULT_CLOSE_TIMEOUT = 30;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The number of records remaining from the previous poll.
|
* The number of records remaining from the previous poll.
|
||||||
* @since 3.2
|
* @since 3.2
|
||||||
@@ -135,6 +137,8 @@ public class KafkaMessageSource<K, V> extends AbstractMessageSource<Object> impl
|
|||||||
|
|
||||||
private boolean running;
|
private boolean running;
|
||||||
|
|
||||||
|
private Duration closeTimeout = Duration.ofSeconds(DEFAULT_CLOSE_TIMEOUT);
|
||||||
|
|
||||||
private volatile Consumer<K, V> consumer;
|
private volatile Consumer<K, V> consumer;
|
||||||
|
|
||||||
private volatile boolean pausing;
|
private volatile boolean pausing;
|
||||||
@@ -334,6 +338,15 @@ public class KafkaMessageSource<K, V> extends AbstractMessageSource<Object> impl
|
|||||||
return this.commitTimeout;
|
return this.commitTimeout;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Set the close timeout - default 30 seconds.
|
||||||
|
* @param closeTimeout the close timeout.
|
||||||
|
*/
|
||||||
|
public void setCloseTimeout(Duration closeTimeout) {
|
||||||
|
Assert.notNull(closeTimeout, "'closeTimeout' cannot be null");
|
||||||
|
this.closeTimeout = closeTimeout;
|
||||||
|
}
|
||||||
|
|
||||||
private ConsumerFactory<K, V> fixOrRejectConsumerFactory(ConsumerFactory<K, V> suppliedConsumerFactory,
|
private ConsumerFactory<K, V> fixOrRejectConsumerFactory(ConsumerFactory<K, V> suppliedConsumerFactory,
|
||||||
boolean allowMultiFetch) {
|
boolean allowMultiFetch) {
|
||||||
|
|
||||||
@@ -571,7 +584,7 @@ public class KafkaMessageSource<K, V> extends AbstractMessageSource<Object> impl
|
|||||||
private void stopConsumer() {
|
private void stopConsumer() {
|
||||||
synchronized (this.consumerMonitor) {
|
synchronized (this.consumerMonitor) {
|
||||||
if (this.consumer != null) {
|
if (this.consumer != null) {
|
||||||
this.consumer.close();
|
this.consumer.close(this.closeTimeout);
|
||||||
this.consumer = null;
|
this.consumer = null;
|
||||||
this.assignedPartitions.clear();
|
this.assignedPartitions.clear();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -21,7 +21,6 @@ import static org.assertj.core.api.Assertions.assertThatIllegalArgumentException
|
|||||||
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
import static org.assertj.core.api.Assertions.assertThatThrownBy;
|
||||||
import static org.mockito.ArgumentMatchers.any;
|
import static org.mockito.ArgumentMatchers.any;
|
||||||
import static org.mockito.ArgumentMatchers.anyCollection;
|
import static org.mockito.ArgumentMatchers.anyCollection;
|
||||||
import static org.mockito.ArgumentMatchers.anyLong;
|
|
||||||
import static org.mockito.ArgumentMatchers.anyString;
|
import static org.mockito.ArgumentMatchers.anyString;
|
||||||
import static org.mockito.ArgumentMatchers.isNull;
|
import static org.mockito.ArgumentMatchers.isNull;
|
||||||
import static org.mockito.BDDMockito.given;
|
import static org.mockito.BDDMockito.given;
|
||||||
@@ -45,7 +44,6 @@ import java.util.LinkedHashSet;
|
|||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.concurrent.TimeUnit;
|
|
||||||
import java.util.concurrent.atomic.AtomicBoolean;
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
import java.util.concurrent.atomic.AtomicInteger;
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
import java.util.concurrent.atomic.AtomicReference;
|
import java.util.concurrent.atomic.AtomicReference;
|
||||||
@@ -329,7 +327,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).resume(partitions.getAllValues().get(1));
|
inOrder.verify(consumer).resume(partitions.getAllValues().get(1));
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verifyNoMoreInteractions();
|
inOrder.verifyNoMoreInteractions();
|
||||||
if (!sync) {
|
if (!sync) {
|
||||||
assertThat(callbackCount.get()).isEqualTo(4);
|
assertThat(callbackCount.get()).isEqualTo(4);
|
||||||
@@ -444,7 +442,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L)));
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L)));
|
||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(6L)));
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(6L)));
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verifyNoMoreInteractions();
|
inOrder.verifyNoMoreInteractions();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -513,7 +511,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L)),
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L)),
|
||||||
Duration.ofSeconds(30));
|
Duration.ofSeconds(30));
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verifyNoMoreInteractions();
|
inOrder.verifyNoMoreInteractions();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -599,7 +597,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L)));
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(2L)));
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verifyNoMoreInteractions();
|
inOrder.verifyNoMoreInteractions();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -744,7 +742,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L)));
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(3L)));
|
||||||
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(4L)));
|
inOrder.verify(consumer).commitSync(Collections.singletonMap(topicPartition, new OffsetAndMetadata(4L)));
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verifyNoMoreInteractions();
|
inOrder.verifyNoMoreInteractions();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -922,9 +920,7 @@ class MessageSourceTests {
|
|||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).resume(anyCollection());
|
inOrder.verify(consumer).resume(anyCollection());
|
||||||
inOrder.verify(consumer).poll(any(Duration.class));
|
inOrder.verify(consumer).poll(any(Duration.class));
|
||||||
inOrder.verify(consumer).close();
|
inOrder.verify(consumer).close(any());
|
||||||
inOrder.verify(consumer).close(anyLong(), any(TimeUnit.class));
|
|
||||||
inOrder.verifyNoMoreInteractions();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -48,6 +48,7 @@ import java.util.concurrent.atomic.AtomicReference;
|
|||||||
|
|
||||||
import org.apache.kafka.clients.consumer.Consumer;
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.apache.kafka.clients.consumer.ConsumerConfig;
|
import org.apache.kafka.clients.consumer.ConsumerConfig;
|
||||||
|
import org.apache.kafka.clients.consumer.ConsumerGroupMetadata;
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
|
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||||
import org.apache.kafka.clients.consumer.ConsumerRecords;
|
import org.apache.kafka.clients.consumer.ConsumerRecords;
|
||||||
@@ -505,6 +506,8 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
}).given(mockConsumer).poll(any(Duration.class));
|
}).given(mockConsumer).poll(any(Duration.class));
|
||||||
|
ConsumerGroupMetadata meta = new ConsumerGroupMetadata("group");
|
||||||
|
given(mockConsumer.groupMetadata()).willReturn(meta);
|
||||||
ConsumerFactory cf = mock(ConsumerFactory.class);
|
ConsumerFactory cf = mock(ConsumerFactory.class);
|
||||||
willReturn(mockConsumer).given(cf).createConsumer("group", "", null, KafkaTestUtils.defaultPropertyOverrides());
|
willReturn(mockConsumer).given(cf).createConsumer("group", "", null, KafkaTestUtils.defaultPropertyOverrides());
|
||||||
Producer producer = mock(Producer.class);
|
Producer producer = mock(Producer.class);
|
||||||
@@ -543,7 +546,7 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
InOrder inOrder = inOrder(producer);
|
InOrder inOrder = inOrder(producer);
|
||||||
inOrder.verify(producer).beginTransaction();
|
inOrder.verify(producer).beginTransaction();
|
||||||
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
||||||
new OffsetAndMetadata(0)), "group");
|
new OffsetAndMetadata(0)), meta);
|
||||||
inOrder.verify(producer).commitTransaction();
|
inOrder.verify(producer).commitTransaction();
|
||||||
inOrder.verify(producer).close(any());
|
inOrder.verify(producer).close(any());
|
||||||
inOrder.verify(producer).beginTransaction();
|
inOrder.verify(producer).beginTransaction();
|
||||||
@@ -551,7 +554,7 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
inOrder.verify(producer).send(captor.capture(), any(Callback.class));
|
inOrder.verify(producer).send(captor.capture(), any(Callback.class));
|
||||||
assertThat(captor.getValue()).isEqualTo(new ProducerRecord("topic", null, "bar", "value"));
|
assertThat(captor.getValue()).isEqualTo(new ProducerRecord("topic", null, "bar", "value"));
|
||||||
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
||||||
new OffsetAndMetadata(1)), "group");
|
new OffsetAndMetadata(1)), meta);
|
||||||
inOrder.verify(producer).commitTransaction();
|
inOrder.verify(producer).commitTransaction();
|
||||||
inOrder.verify(producer).close(any());
|
inOrder.verify(producer).close(any());
|
||||||
container.stop();
|
container.stop();
|
||||||
@@ -617,6 +620,8 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
}).given(mockConsumer).poll(any(Duration.class));
|
}).given(mockConsumer).poll(any(Duration.class));
|
||||||
|
ConsumerGroupMetadata meta = new ConsumerGroupMetadata("group");
|
||||||
|
given(mockConsumer.groupMetadata()).willReturn(meta);
|
||||||
ConsumerFactory cf = mock(ConsumerFactory.class);
|
ConsumerFactory cf = mock(ConsumerFactory.class);
|
||||||
willReturn(mockConsumer).given(cf).createConsumer("group", "", null, KafkaTestUtils.defaultPropertyOverrides());
|
willReturn(mockConsumer).given(cf).createConsumer("group", "", null, KafkaTestUtils.defaultPropertyOverrides());
|
||||||
Producer producer = mock(Producer.class);
|
Producer producer = mock(Producer.class);
|
||||||
@@ -663,7 +668,7 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
InOrder inOrder = inOrder(producer);
|
InOrder inOrder = inOrder(producer);
|
||||||
inOrder.verify(producer).beginTransaction();
|
inOrder.verify(producer).beginTransaction();
|
||||||
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
||||||
new OffsetAndMetadata(0)), "group");
|
new OffsetAndMetadata(0)), meta);
|
||||||
inOrder.verify(producer).commitTransaction();
|
inOrder.verify(producer).commitTransaction();
|
||||||
inOrder.verify(producer).close(any());
|
inOrder.verify(producer).close(any());
|
||||||
inOrder.verify(producer).beginTransaction();
|
inOrder.verify(producer).beginTransaction();
|
||||||
@@ -671,7 +676,7 @@ class KafkaProducerMessageHandlerTests {
|
|||||||
inOrder.verify(producer).send(captor.capture(), any(Callback.class));
|
inOrder.verify(producer).send(captor.capture(), any(Callback.class));
|
||||||
assertThat(captor.getValue()).isEqualTo(new ProducerRecord("topic", null, "bar", "value"));
|
assertThat(captor.getValue()).isEqualTo(new ProducerRecord("topic", null, "bar", "value"));
|
||||||
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
inOrder.verify(producer).sendOffsetsToTransaction(Collections.singletonMap(topicPartition,
|
||||||
new OffsetAndMetadata(1)), "group");
|
new OffsetAndMetadata(1)), meta);
|
||||||
inOrder.verify(producer).commitTransaction();
|
inOrder.verify(producer).commitTransaction();
|
||||||
inOrder.verify(producer).close(any());
|
inOrder.verify(producer).close(any());
|
||||||
container.stop();
|
container.stop();
|
||||||
|
|||||||
Reference in New Issue
Block a user