From 7974331b338546affff197c28be0150b4462c164 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Wed, 17 Jun 2020 12:57:41 -0400 Subject: [PATCH] GH-1508: Fix DefaultARProcessor for EOSMode.BETA Resolves https://github.com/spring-projects/spring-kafka/issues/1508 - also validate configuration when `commitRecovered` is `true`; deprecate setters. --- .../listener/AfterRollbackProcessor.java | 42 +++++++-- .../DefaultAfterRollbackProcessor.java | 94 +++++++++++++++---- .../KafkaMessageListenerContainer.java | 10 +- .../DefaultAfterRollbackProcessorTests.java | 34 ++++--- .../listener/TransactionalContainerTests.java | 11 +-- 5 files changed, 142 insertions(+), 49 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AfterRollbackProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AfterRollbackProcessor.java index 28202fba..7ba4abfc 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AfterRollbackProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AfterRollbackProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2018-2019 the original author or authors. + * Copyright 2018-2020 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. @@ -21,6 +21,8 @@ import java.util.List; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.springframework.kafka.listener.ContainerProperties.EOSMode; + /** * Invoked by a listener container with remaining, unprocessed, records * (including the failed record). Implementations should seek the desired @@ -55,9 +57,37 @@ public interface AfterRollbackProcessor { * @param recoverable the recoverable. * @since 2.2 * @see #isProcessInTransaction() + * @deprecated in favor of {@link #process(List, Consumer, Exception, boolean, + * EOSMode)}. */ + @Deprecated void process(List> records, Consumer consumer, Exception exception, boolean recoverable); + /** + * Process the remaining records. Recoverable will be true if the container is + * processing individual records; this allows the processor to recover (skip) the + * failed record rather than re-seeking it. This is not possible with a batch listener + * since only the listener itself knows which record in the batch keeps failing. + * IMPORTANT: If invoked in a transaction when the listener was invoked with a single + * record, the transaction id will be based on the container group.id and the + * topic/partition of the failed record, to avoid issues with zombie fencing (unless + * the {@link EOSMode} is {@link EOSMode#BETA}). So, generally, only its offset should + * be sent to the transaction. For other behavior the process method should manage its + * own transaction. + * @param records the records. + * @param consumer the consumer. + * @param exception the exception + * @param recoverable the recoverable. + * @param eosMode the {@link EOSMode}. + * @since 2.5.3 + * @see #isProcessInTransaction() + */ + default void process(List> records, Consumer consumer, Exception exception, + boolean recoverable, EOSMode eosMode) { + + process(records, consumer, exception, recoverable); + } + /** * Optional method to clear thread state; will be called just before a consumer * thread terminates. @@ -68,13 +98,13 @@ public interface AfterRollbackProcessor { } /** - * Return true to invoke {@link #process(List, Consumer, Exception, boolean)} in a new - * transaction. Because the container cannot infer the desired behavior, the processor - * is responsible for sending the offset to the transaction if it decides to skip the - * failing record. + * Return true to invoke {@link #process(List, Consumer, Exception, boolean, EOSMode)} + * in a new transaction. Because the container cannot infer the desired behavior, the + * processor is responsible for sending the offset to the transaction if it decides to + * skip the failing record. * @return true to run in a transaction; default false. * @since 2.2.5 - * @see #process(List, Consumer, Exception, boolean) + * @see #process(List, Consumer, Exception, boolean, EOSMode) */ default boolean isProcessInTransaction() { return false; diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessor.java index 40ca1592..9d367b28 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessor.java @@ -25,9 +25,12 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; +import org.springframework.beans.factory.InitializingBean; import org.springframework.kafka.core.KafkaOperations; import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.listener.ContainerProperties.EOSMode; import org.springframework.lang.Nullable; +import org.springframework.util.Assert; import org.springframework.util.backoff.BackOff; /** @@ -46,9 +49,10 @@ import org.springframework.util.backoff.BackOff; * @since 1.3.5 * */ -public class DefaultAfterRollbackProcessor extends FailedRecordProcessor implements AfterRollbackProcessor { +public class DefaultAfterRollbackProcessor extends FailedRecordProcessor + implements AfterRollbackProcessor, InitializingBean { - private KafkaOperations kafkaTemplate; + private KafkaOperations kafkaTemplate; /** * Construct an instance with the default recoverer which simply logs the record after @@ -91,21 +95,66 @@ public class DefaultAfterRollbackProcessor extends FailedRecordProcessor i public DefaultAfterRollbackProcessor(@Nullable BiConsumer, Exception> recoverer, BackOff backOff) { + this(recoverer, backOff, null, false); + } + + /** + * Construct an instance with the provided recoverer which will be called after the + * backOff returns STOP for a topic/partition/offset. + * @param recoverer the recoverer; if null, the default (logging) recoverer is used. + * @param backOff the {@link BackOff}. + * @param kafkaOperations for sending the recovered offset to the transaction. + * @param commitRecovered true to commit the recovered record's offset; requires a + * {@link KafkaOperations}. + * @since 2.5.3 + */ + public DefaultAfterRollbackProcessor(@Nullable BiConsumer, Exception> recoverer, + BackOff backOff, @Nullable KafkaOperations kafkaOperations, boolean commitRecovered) { + super(recoverer, backOff); + this.kafkaTemplate = kafkaOperations; + super.setCommitRecovered(commitRecovered); + checkConfig(); + } + + @Override + public void afterPropertiesSet() { + // remove InitializingBean when the deprecated setters are removed. + checkConfig(); + } + + private void checkConfig() { + Assert.isTrue(!isCommitRecovered() || this.kafkaTemplate != null, + "A KafkaOperations is required when 'commitRecovered' is true"); + } + + @Override + @Deprecated + public void process(List> records, Consumer consumer, Exception exception, + boolean recoverable) { + + process(records, consumer, exception, recoverable, EOSMode.ALPHA); } @SuppressWarnings({ "unchecked", "rawtypes" }) @Override public void process(List> records, Consumer consumer, Exception exception, - boolean recoverable) { + boolean recoverable, @Nullable EOSMode eosMode) { if (SeekUtils.doSeeks(((List) records), consumer, exception, recoverable, getSkipPredicate((List) records, exception), this.logger) - && isCommitRecovered() && this.kafkaTemplate != null && this.kafkaTemplate.isTransactional()) { + && isCommitRecovered() && this.kafkaTemplate.isTransactional()) { ConsumerRecord skipped = records.get(0); - this.kafkaTemplate.sendOffsetsToTransaction( - Collections.singletonMap(new TopicPartition(skipped.topic(), skipped.partition()), - new OffsetAndMetadata(skipped.offset() + 1))); + if (EOSMode.ALPHA.equals(eosMode)) { + this.kafkaTemplate.sendOffsetsToTransaction( + Collections.singletonMap(new TopicPartition(skipped.topic(), skipped.partition()), + new OffsetAndMetadata(skipped.offset() + 1))); + } + else { + this.kafkaTemplate.sendOffsetsToTransaction( + Collections.singletonMap(new TopicPartition(skipped.topic(), skipped.partition()), + new OffsetAndMetadata(skipped.offset() + 1)), consumer.groupMetadata()); + } } } @@ -115,29 +164,32 @@ public class DefaultAfterRollbackProcessor extends FailedRecordProcessor i } /** - * {@inheritDoc} - * Set to true and the container will run the - * {@link #process(List, Consumer, Exception, boolean)} method in a transaction and, - * if a record is skipped and recovered, we will send its offset to the transaction. - * Requires a {@link KafkaOperations}. + * {@inheritDoc} Set to true and the container will run the {@link #process(List, + * Consumer, Exception, boolean, EOSMode)} method in a transaction and, if a record + * is skipped and recovered, we will send its offset to the transaction. Requires a + * {@link KafkaOperations}. * @param commitRecovered true to process in a transaction. * @since 2.3 + * @deprecated in favor of + * {@link #DefaultAfterRollbackProcessor(BiConsumer, BackOff, KafkaOperations, boolean)}. * @see #isProcessInTransaction() - * @see #process(List, Consumer, Exception, boolean) - * @see #setKafkaOperations(KafkaOperations) + * @see #process(List, Consumer, Exception, boolean, EOSMode) */ + @Deprecated @Override public void setCommitRecovered(boolean commitRecovered) { // NOSONAR enhanced javadoc super.setCommitRecovered(commitRecovered); } /** - * Set a {@link KafkaTemplate} to use to send the offset of a recovered record - * to a transaction. + * Set a {@link KafkaTemplate} to use to send the offset of a recovered record to a + * transaction. * @param kafkaTemplate the template. * @since 2.2.5 + * @deprecated in favor of + * {@link #DefaultAfterRollbackProcessor(BiConsumer, BackOff, KafkaOperations, + * boolean)}. * @see #setCommitRecovered(boolean) - * @deprecated in favor of {@link #setKafkaOperations(KafkaOperations)}. */ @Deprecated public void setKafkaTemplate(KafkaTemplate kafkaTemplate) { @@ -145,12 +197,16 @@ public class DefaultAfterRollbackProcessor extends FailedRecordProcessor i } /** - * Set a {@link KafkaOperations} to use to send the offset of a recovered record - * to a transaction. + * Set a {@link KafkaOperations} to use to send the offset of a recovered record to a + * transaction. * @param kafkaOperations the operations. * @since 2.5.1 + * @deprecated in favor of + * {@link #DefaultAfterRollbackProcessor(BiConsumer, BackOff, KafkaOperations, + * boolean)}. * @see #setCommitRecovered(boolean) */ + @Deprecated public void setKafkaOperations(KafkaOperations kafkaOperations) { this.kafkaTemplate = kafkaOperations; } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 22d51e16..552c940a 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -1430,10 +1430,11 @@ public class KafkaMessageListenerContainer // NOSONAR line count try { if (recordList == null) { afterRollbackProcessorToUse.process(createRecordList(records), this.consumer, rollbackException, - false); + false, this.eosMode); } else { - afterRollbackProcessorToUse.process(recordList, this.consumer, rollbackException, false); + afterRollbackProcessorToUse.process(recordList, this.consumer, rollbackException, false, + this.eosMode); } } catch (KafkaException ke) { @@ -1690,14 +1691,15 @@ public class KafkaMessageListenerContainer // NOSONAR line count @Override protected void doInTransactionWithoutResult(TransactionStatus status) { - afterRollbackProcessorToUse.process(unprocessed, ListenerConsumer.this.consumer, e, true); + afterRollbackProcessorToUse.process(unprocessed, ListenerConsumer.this.consumer, e, true, + ListenerConsumer.this.eosMode); } }); } else { try { - afterRollbackProcessorToUse.process(unprocessed, this.consumer, e, true); + afterRollbackProcessorToUse.process(unprocessed, this.consumer, e, true, this.eosMode); } catch (KafkaException ke) { ke.selfLog("AfterRollbackProcessor threw an exception", this.logger); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessorTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessorTests.java index 0c9e58f1..c86d67b0 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessorTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/DefaultAfterRollbackProcessorTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2019 the original author or authors. + * Copyright 2019-2020 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. @@ -17,10 +17,12 @@ package org.springframework.kafka.listener; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.BDDMockito.given; import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -30,12 +32,14 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerGroupMetadata; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.api.Test; import org.mockito.InOrder; -import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.kafka.core.KafkaOperations; +import org.springframework.kafka.listener.ContainerProperties.EOSMode; import org.springframework.kafka.support.serializer.DeserializationException; /** @@ -49,40 +53,42 @@ public class DefaultAfterRollbackProcessorTests { public void testClassifier() { AtomicReference> recovered = new AtomicReference<>(); AtomicBoolean recovererShouldFail = new AtomicBoolean(false); + @SuppressWarnings("unchecked") + KafkaOperations template = mock(KafkaOperations.class); + given(template.isTransactional()).willReturn(true); DefaultAfterRollbackProcessor processor = new DefaultAfterRollbackProcessor<>((r, t) -> { if (recovererShouldFail.getAndSet(false)) { throw new RuntimeException("test recoverer failure"); } recovered.set(r); - }); - @SuppressWarnings("unchecked") - KafkaTemplate template = mock(KafkaTemplate.class); - given(template.isTransactional()).willReturn(true); - processor.setKafkaOperations(template); - processor.setCommitRecovered(true); + }, SeekUtils.DEFAULT_BACK_OFF, template, true); ConsumerRecord record1 = new ConsumerRecord<>("foo", 0, 0L, "foo", "bar"); ConsumerRecord record2 = new ConsumerRecord<>("foo", 1, 1L, "foo", "bar"); List> records = Arrays.asList(record1, record2); IllegalStateException illegalState = new IllegalStateException(); @SuppressWarnings("unchecked") Consumer consumer = mock(Consumer.class); - processor.process(records, consumer, illegalState, true); - processor.process(records, consumer, new DeserializationException("intended", null, false, illegalState), true); + given(consumer.groupMetadata()).willReturn(new ConsumerGroupMetadata("foo")); + processor.process(records, consumer, illegalState, true, EOSMode.ALPHA); + processor.process(records, consumer, new DeserializationException("intended", null, false, illegalState), true, + EOSMode.ALPHA); verify(template).sendOffsetsToTransaction(anyMap()); + verify(template, never()).sendOffsetsToTransaction(anyMap(), any(ConsumerGroupMetadata.class)); assertThat(recovered.get()).isSameAs(record1); processor.addNotRetryableException(IllegalStateException.class); recovered.set(null); recovererShouldFail.set(true); - processor.process(records, consumer, illegalState, true); - verify(template).sendOffsetsToTransaction(anyMap()); // recovery failed - processor.process(records, consumer, illegalState, true); - verify(template, times(2)).sendOffsetsToTransaction(anyMap()); + processor.process(records, consumer, illegalState, true, EOSMode.ALPHA); + verify(template, times(1)).sendOffsetsToTransaction(anyMap()); // recovery failed + processor.process(records, consumer, illegalState, true, EOSMode.BETA); + verify(template, times(1)).sendOffsetsToTransaction(anyMap(), any(ConsumerGroupMetadata.class)); assertThat(recovered.get()).isSameAs(record1); InOrder inOrder = inOrder(consumer); inOrder.verify(consumer).seek(new TopicPartition("foo", 0), 0L); // not recovered so seek inOrder.verify(consumer, times(2)).seek(new TopicPartition("foo", 1), 1L); inOrder.verify(consumer).seek(new TopicPartition("foo", 0), 0L); // recovery failed inOrder.verify(consumer, times(2)).seek(new TopicPartition("foo", 1), 1L); + inOrder.verify(consumer).groupMetadata(); inOrder.verifyNoMoreInteractions(); } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/TransactionalContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/TransactionalContainerTests.java index fa04796a..2e4efcaf 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/TransactionalContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/TransactionalContainerTests.java @@ -571,7 +571,7 @@ public class TransactionalContainerTests { consumer.close(); } - @SuppressWarnings("unchecked") + @SuppressWarnings({ "unchecked", "deprecation" }) @Test public void testMaxFailures() throws Exception { logger.info("Start testMaxFailures"); @@ -619,9 +619,7 @@ public class TransactionalContainerTests { }; DefaultAfterRollbackProcessor afterRollbackProcessor = - spy(new DefaultAfterRollbackProcessor<>(recoverer, new FixedBackOff(0L, 2L))); - afterRollbackProcessor.setCommitRecovered(true); - afterRollbackProcessor.setKafkaOperations(dlTemplate); + spy(new DefaultAfterRollbackProcessor<>(recoverer, new FixedBackOff(0L, 2L), dlTemplate, true)); container.setAfterRollbackProcessor(afterRollbackProcessor); final CountDownLatch stopLatch = new CountDownLatch(1); container.setApplicationEventPublisher(e -> { @@ -666,7 +664,8 @@ public class TransactionalContainerTests { assertThat(stopLatch.await(10, TimeUnit.SECONDS)).isTrue(); verify(afterRollbackProcessor, times(4)).isProcessInTransaction(); ArgumentCaptor captor = ArgumentCaptor.forClass(Exception.class); - verify(afterRollbackProcessor, times(4)).process(any(), any(), captor.capture(), anyBoolean()); + verify(afterRollbackProcessor, never()).process(any(), any(), captor.capture(), anyBoolean()); + verify(afterRollbackProcessor, times(4)).process(any(), any(), captor.capture(), anyBoolean(), any()); assertThat(captor.getValue()).isInstanceOf(ListenerExecutionFailedException.class) .extracting(ex -> ((ListenerExecutionFailedException) ex).getGroupId()) .isEqualTo("groupInARBP"); @@ -794,7 +793,7 @@ public class TransactionalContainerTests { inOrder.verifyNoMoreInteractions(); assertThat(deliveryCount.get()).isEqualTo(1); - verify(arp, never()).process(any(), any(), any(), anyBoolean()); + verify(arp, never()).process(any(), any(), any(), anyBoolean(), any()); container.stop(); }