Updates for latest S-K snapshot
- change `TopicPartitionInitialOffset` to `TopicPartitionOffset` - when resetting with manual assignment, set the SeekPosition earlier rather than relying on direct updat of the `ContainerProperties` field - set `missingTopicsFatal` to `false` in mock test to avoid log messages about missing bootstrap servers
This commit is contained in:
@@ -100,8 +100,8 @@ import org.springframework.kafka.support.KafkaHeaderMapper;
|
||||
import org.springframework.kafka.support.KafkaHeaders;
|
||||
import org.springframework.kafka.support.ProducerListener;
|
||||
import org.springframework.kafka.support.SendResult;
|
||||
import org.springframework.kafka.support.TopicPartitionInitialOffset;
|
||||
import org.springframework.kafka.support.TopicPartitionInitialOffset.SeekPosition;
|
||||
import org.springframework.kafka.support.TopicPartitionOffset;
|
||||
import org.springframework.kafka.support.TopicPartitionOffset.SeekPosition;
|
||||
import org.springframework.kafka.support.converter.MessagingMessageConverter;
|
||||
import org.springframework.kafka.transaction.KafkaTransactionManager;
|
||||
import org.springframework.lang.Nullable;
|
||||
@@ -510,14 +510,15 @@ public class KafkaMessageChannelBinder extends
|
||||
Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions),
|
||||
"A list of partitions must be provided");
|
||||
}
|
||||
final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets(
|
||||
listenedPartitions);
|
||||
final TopicPartitionOffset[] topicPartitionOffsets = groupManagement
|
||||
? null
|
||||
: getTopicPartitionOffsets(listenedPartitions, extendedConsumerProperties, consumerFactory);
|
||||
final ContainerProperties containerProperties = anonymous
|
||||
|| extendedConsumerProperties.getExtension().isAutoRebalanceEnabled()
|
||||
|| groupManagement
|
||||
? usingPatterns
|
||||
? new ContainerProperties(Pattern.compile(topics[0]))
|
||||
: new ContainerProperties(topics)
|
||||
: new ContainerProperties(topicPartitionInitialOffsets);
|
||||
: new ContainerProperties(topicPartitionOffsets);
|
||||
if (this.transactionManager != null) {
|
||||
containerProperties.setTransactionManager(this.transactionManager);
|
||||
}
|
||||
@@ -534,8 +535,7 @@ public class KafkaMessageChannelBinder extends
|
||||
if (groupManagement && listenedPartitions.isEmpty()) {
|
||||
concurrency = extendedConsumerProperties.getConcurrency();
|
||||
}
|
||||
resetOffsets(extendedConsumerProperties, consumerFactory, groupManagement,
|
||||
containerProperties);
|
||||
resetOffsetsForAutoRebalance(extendedConsumerProperties, consumerFactory, containerProperties);
|
||||
@SuppressWarnings("rawtypes")
|
||||
final ConcurrentMessageListenerContainer<?, ?> messageListenerContainer = new ConcurrentMessageListenerContainer(
|
||||
consumerFactory, containerProperties) {
|
||||
@@ -681,20 +681,14 @@ public class KafkaMessageChannelBinder extends
|
||||
* Reset the offsets if needed; may update the offsets in in the container's
|
||||
* topicPartitionInitialOffsets.
|
||||
*/
|
||||
private void resetOffsets(
|
||||
private void resetOffsetsForAutoRebalance(
|
||||
final ExtendedConsumerProperties<KafkaConsumerProperties> extendedConsumerProperties,
|
||||
final ConsumerFactory<?, ?> consumerFactory, boolean groupManagement,
|
||||
final ContainerProperties containerProperties) {
|
||||
final ConsumerFactory<?, ?> consumerFactory, final ContainerProperties containerProperties) {
|
||||
|
||||
boolean resetOffsets = extendedConsumerProperties.getExtension().isResetOffsets();
|
||||
final Object resetTo = consumerFactory.getConfigurationProperties()
|
||||
.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG);
|
||||
if (!"earliest".equals(resetTo) && !"latest".equals(resetTo)) {
|
||||
logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG
|
||||
+ " property cannot reset");
|
||||
resetOffsets = false;
|
||||
}
|
||||
if (groupManagement && resetOffsets) {
|
||||
final Object resetTo = checkReset(extendedConsumerProperties.getExtension().isResetOffsets(),
|
||||
consumerFactory.getConfigurationProperties()
|
||||
.get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
|
||||
if (resetTo != null) {
|
||||
Set<TopicPartition> sought = ConcurrentHashMap.newKeySet();
|
||||
containerProperties.setConsumerRebalanceListener(new ConsumerAwareRebalanceListener() {
|
||||
|
||||
@@ -736,15 +730,15 @@ public class KafkaMessageChannelBinder extends
|
||||
}
|
||||
});
|
||||
}
|
||||
else if (resetOffsets) {
|
||||
Arrays.stream(containerProperties.getTopicPartitions())
|
||||
.map(tpio -> new TopicPartitionInitialOffset(tpio.topic(),
|
||||
tpio.partition(),
|
||||
"earliest".equals(resetTo) ? SeekPosition.BEGINNING
|
||||
: SeekPosition.END))
|
||||
.collect(Collectors.toList())
|
||||
.toArray(containerProperties.getTopicPartitions());
|
||||
}
|
||||
|
||||
private Object checkReset(boolean resetOffsets, final Object resetTo) {
|
||||
if (resetOffsets && !"earliest".equals(resetTo) && !"latest".equals(resetTo)) {
|
||||
logger.warn("no (or unknown) " + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG
|
||||
+ " property cannot reset");
|
||||
return null;
|
||||
}
|
||||
return resetTo;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -1118,17 +1112,26 @@ public class KafkaMessageChannelBinder extends
|
||||
&& properties.getExtension().isEnableDlq();
|
||||
}
|
||||
|
||||
private TopicPartitionInitialOffset[] getTopicPartitionInitialOffsets(
|
||||
Collection<PartitionInfo> listenedPartitions) {
|
||||
final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = new TopicPartitionInitialOffset[listenedPartitions
|
||||
.size()];
|
||||
private TopicPartitionOffset[] getTopicPartitionOffsets(
|
||||
Collection<PartitionInfo> listenedPartitions,
|
||||
ExtendedConsumerProperties<KafkaConsumerProperties> extendedConsumerProperties,
|
||||
ConsumerFactory<?, ?> consumerFactory) {
|
||||
|
||||
final TopicPartitionOffset[] TopicPartitionOffsets =
|
||||
new TopicPartitionOffset[listenedPartitions.size()];
|
||||
int i = 0;
|
||||
SeekPosition seekPosition = null;
|
||||
Object resetTo = checkReset(extendedConsumerProperties.getExtension().isResetOffsets(),
|
||||
consumerFactory.getConfigurationProperties().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG));
|
||||
if (resetTo != null) {
|
||||
seekPosition = "earliest".equals(resetTo) ? SeekPosition.BEGINNING : SeekPosition.END;
|
||||
}
|
||||
for (PartitionInfo partition : listenedPartitions) {
|
||||
|
||||
topicPartitionInitialOffsets[i++] = new TopicPartitionInitialOffset(
|
||||
partition.topic(), partition.partition());
|
||||
TopicPartitionOffsets[i++] = new TopicPartitionOffset(
|
||||
partition.topic(), partition.partition(), seekPosition);
|
||||
}
|
||||
return topicPartitionInitialOffsets;
|
||||
return TopicPartitionOffsets;
|
||||
}
|
||||
|
||||
private String toDisplayString(String original, int maxCharacters) {
|
||||
|
||||
@@ -114,7 +114,7 @@ import org.springframework.kafka.listener.MessageListenerContainer;
|
||||
import org.springframework.kafka.support.Acknowledgment;
|
||||
import org.springframework.kafka.support.KafkaHeaders;
|
||||
import org.springframework.kafka.support.SendResult;
|
||||
import org.springframework.kafka.support.TopicPartitionInitialOffset;
|
||||
import org.springframework.kafka.support.TopicPartitionOffset;
|
||||
import org.springframework.kafka.support.converter.MessagingMessageConverter;
|
||||
import org.springframework.kafka.test.core.BrokerAddress;
|
||||
import org.springframework.kafka.test.rule.EmbeddedKafkaRule;
|
||||
@@ -2455,14 +2455,15 @@ public class KafkaBinderTests extends
|
||||
binding = binder.bindConsumer(testTopicName, "test-x", input,
|
||||
consumerProperties);
|
||||
|
||||
TopicPartitionInitialOffset[] listenedPartitions = TestUtils.getPropertyValue(
|
||||
ContainerProperties containerProps = TestUtils.getPropertyValue(
|
||||
binding,
|
||||
"lifecycle.messageListenerContainer.containerProperties.topicPartitions",
|
||||
TopicPartitionInitialOffset[].class);
|
||||
"lifecycle.messageListenerContainer.containerProperties",
|
||||
ContainerProperties.class);
|
||||
TopicPartitionOffset[] listenedPartitions = containerProps.getTopicPartitionsToAssign();
|
||||
assertThat(listenedPartitions).hasSize(2);
|
||||
assertThat(listenedPartitions).contains(
|
||||
new TopicPartitionInitialOffset(testTopicName, 2),
|
||||
new TopicPartitionInitialOffset(testTopicName, 5));
|
||||
new TopicPartitionOffset(testTopicName, 2),
|
||||
new TopicPartitionOffset(testTopicName, 5));
|
||||
int partitions = invokePartitionSize(testTopicName);
|
||||
assertThat(partitions).isEqualTo(6);
|
||||
}
|
||||
|
||||
@@ -45,11 +45,13 @@ import org.springframework.cloud.stream.binder.ExtendedConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties;
|
||||
import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner;
|
||||
import org.springframework.cloud.stream.config.ListenerContainerCustomizer;
|
||||
import org.springframework.cloud.stream.provisioning.ConsumerDestination;
|
||||
import org.springframework.context.support.GenericApplicationContext;
|
||||
import org.springframework.integration.channel.DirectChannel;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.kafka.core.ConsumerFactory;
|
||||
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
@@ -175,6 +177,7 @@ public class KafkaBinderUnitTests {
|
||||
|
||||
private void testOffsetResetWithGroupManagement(final boolean earliest,
|
||||
boolean groupManage, String topic, String group) throws Exception {
|
||||
|
||||
final List<TopicPartition> partitions = new ArrayList<>();
|
||||
partitions.add(new TopicPartition(topic, 0));
|
||||
partitions.add(new TopicPartition(topic, 1));
|
||||
@@ -218,8 +221,18 @@ public class KafkaBinderUnitTests {
|
||||
latch.countDown();
|
||||
return null;
|
||||
}).given(consumer).seekToEnd(any());
|
||||
class Customizer implements ListenerContainerCustomizer<AbstractMessageListenerContainer<?, ?>> {
|
||||
|
||||
@Override
|
||||
public void configure(AbstractMessageListenerContainer<?, ?> container, String destinationName,
|
||||
String group) {
|
||||
|
||||
container.getContainerProperties().setMissingTopicsFatal(false);
|
||||
}
|
||||
|
||||
}
|
||||
KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(
|
||||
configurationProperties, provisioningProvider) {
|
||||
configurationProperties, provisioningProvider, new Customizer(), null) {
|
||||
|
||||
@Override
|
||||
protected ConsumerFactory<?, ?> createKafkaConsumerFactory(boolean anonymous,
|
||||
|
||||
Reference in New Issue
Block a user