From 3bfd87c37032b1a6bdd343e2ee135d0b3bf33b6f Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Wed, 5 Feb 2020 12:44:05 -0800 Subject: [PATCH 01/17] add timer for update limit offsets --- .../processor/internals/ChangelogReader.java | 5 - .../internals/StoreChangelogReader.java | 103 ++++++++++-------- .../processor/internals/StreamThread.java | 2 +- .../internals/MockChangelogReader.java | 5 - .../internals/StoreChangelogReaderTest.java | 97 ++++++++++------- .../StreamThreadStateStoreProviderTest.java | 1 + .../apache/kafka/test/StreamsTestUtils.java | 6 +- .../kafka/streams/TopologyTestDriver.java | 1 + 8 files changed, 121 insertions(+), 99 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ChangelogReader.java index 2f31ad136626a..7fe93c8bcad54 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ChangelogReader.java @@ -30,11 +30,6 @@ public interface ChangelogReader extends ChangelogRegister { */ void restore(); - /** - * Update offset limit of a given changelog partition - */ - void updateLimitOffsets(); - /** * Transit to restore active changelogs mode */ diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index b813961682306..dac6115c2a399 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -24,6 +24,7 @@ import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.errors.TaskMigratedException; @@ -99,11 +100,18 @@ static class ChangelogMetadata { // only for active restoring tasks (for standby changelog it is null) // NOTE we do not book keep the current offset since we leverage state manager as its source of truth - private Long restoreEndOffset; - // only for standby tasks that use source topics as changelogs (for active it is null); - // if it is not on source topics it is also null - private Long restoreLimitOffset; + + // the end offset beyond which records should not be applied (yet) to restore the states + // + // for both active restoring tasks and standby updating tasks, it is defined as: + // * log-end-offset if the changelog is not piggy-backed with source topic + // * min(log-end-offset, committed-offset) if the changelog is piggy-backed with source topic + // + // the log-end-offset only needs to be updated once and only need to be for active tasks since for standby + // tasks it would never "complete" based on the end-offset; + // the committed-offset needs to be updated periodically for those standby tasks + private Long restoreEndOffset; // buffer records polled by the restore consumer; private final List> bufferedRecords; @@ -113,14 +121,13 @@ static class ChangelogMetadata { private int bufferedLimitIndex; private ChangelogMetadata(final StateStoreMetadata storeMetadata, final ProcessorStateManager stateManager) { + this.changelogState = ChangelogState.REGISTERED; this.storeMetadata = storeMetadata; this.stateManager = stateManager; - this.changelogState = ChangelogState.REGISTERED; this.restoreEndOffset = null; this.totalRestored = 0L; this.bufferedRecords = new ArrayList<>(); - this.restoreLimitOffset = null; this.bufferedLimitIndex = 0; } @@ -139,7 +146,7 @@ private void transitTo(final ChangelogState newState) { public String toString() { final Long currentOffset = storeMetadata.offset(); return changelogState + " " + stateManager.taskType() + - " (currentOffset " + currentOffset + ", endOffset " + restoreEndOffset + ", limitOffset " + restoreLimitOffset + ")"; + " (currentOffset " + currentOffset + ", endOffset " + restoreEndOffset + ")"; } // for testing only below @@ -155,10 +162,6 @@ Long endOffset() { return restoreEndOffset; } - Long limitOffset() { - return restoreLimitOffset; - } - List> bufferedRecords() { return bufferedRecords; } @@ -168,10 +171,14 @@ int bufferedLimitIndex() { } } + private final static long DEFAULT_OFFSET_UPDATE_MS = 5 * 60 * 1000; // five minutes + private ChangelogReaderState state; + private final Time time; private final Logger log; - private final Duration pollTime; + private final Duration pollTimeMs; + private final long updateOffsetIntervalMs; // 1) we keep adding partitions to restore consumer whenever new tasks are registered with the state manager; // 2) we do not unassign partitions when we switch between standbys and actives, we just pause / resume them; @@ -188,18 +195,18 @@ int bufferedLimitIndex() { // to update offset limit for standby tasks; private Consumer mainConsumer; - // the flag indicating limit offsets could be updated --- this is only needed for standby tasks that have limit - // offsets enabled - private boolean updateLimitOffset; + private long lastUpdateOffsetTime; void setMainConsumer(final Consumer consumer) { this.mainConsumer = consumer; } - public StoreChangelogReader(final StreamsConfig config, + public StoreChangelogReader(final Time time, + final StreamsConfig config, final LogContext logContext, final Consumer restoreConsumer, final StateRestoreListener stateRestoreListener) { + this.time = time; this.log = logContext.logger(StoreChangelogReader.class); this.state = ChangelogReaderState.ACTIVE_RESTORING; this.restoreConsumer = restoreConsumer; @@ -208,10 +215,12 @@ public StoreChangelogReader(final StreamsConfig config, // NOTE for restoring active and updating standby we may prefer different poll time // in order to make sure we call the main consumer#poll in time. // TODO: once both of these are moved to a separate thread this may no longer be a concern - this.pollTime = Duration.ofMillis(config.getLong(StreamsConfig.POLL_MS_CONFIG)); + this.pollTimeMs = Duration.ofMillis(config.getLong(StreamsConfig.POLL_MS_CONFIG)); + this.updateOffsetIntervalMs = config.getLong(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG) == Long.MAX_VALUE ? + DEFAULT_OFFSET_UPDATE_MS : config.getLong(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG); + this.lastUpdateOffsetTime = 0L; this.changelogs = new HashMap<>(); - this.updateLimitOffset = false; } private static String recordEndOffset(final Long endOffset) { @@ -311,7 +320,7 @@ public void register(final TopicPartition partition, final ProcessorStateManager // initializing limit offset to 0L for standby changelog to effectively disable any restoration until it is updated if (stateManager.taskType() == Task.TaskType.STANDBY && stateManager.changelogAsSource(partition)) { - changelogMetadata.restoreLimitOffset = 0L; + changelogMetadata.restoreEndOffset = 0L; } if (changelogs.putIfAbsent(partition, changelogMetadata) != null) { @@ -391,16 +400,12 @@ public void restore() { return; } - if (updateLimitOffset) { - updateLimitOffsets(); - } - final Set restoringChangelogs = restoringChangelogs(); if (!restoringChangelogs.isEmpty()) { final ConsumerRecords polledRecords; try { - polledRecords = restoreConsumer.poll(pollTime); + polledRecords = restoreConsumer.poll(pollTimeMs); } catch (final FencedInstanceIdException e) { // when the consumer gets fenced, all its tasks should be migrated throw new TaskMigratedException("Restore consumer get fenced by instance-id polling records.", e); @@ -420,28 +425,26 @@ public void restore() { // small batches; this can be optimized in the future, e.g. wait longer for larger batches. restoreChangelog(changelogs.get(partition)); } - } - // for standby changelogs, if there are buffered records not applicable, it means that the limit offset - // is there to prevent so, we can try to update the limit offset next time. - final Set standbyChangelogs = changelogs.values().stream() - .filter(metadata -> metadata.stateManager.taskType() == Task.TaskType.STANDBY) - .collect(Collectors.toSet()); - for (final ChangelogMetadata metadata : standbyChangelogs) { - if (!metadata.bufferedRecords().isEmpty()) { - updateLimitOffset = true; - break; + + // for standby changelogs, if the interval has elapsed and there are buffered records not applicable, + // we can try to update the limit offset next time. + if (updateOffsetIntervalMs < time.milliseconds() - lastUpdateOffsetTime) { + final Set standbyChangelogs = changelogs.values().stream() + .filter(metadata -> metadata.stateManager.taskType() == Task.TaskType.STANDBY) + .collect(Collectors.toSet()); + for (final ChangelogMetadata metadata : standbyChangelogs) { + if (!metadata.bufferedRecords().isEmpty()) { + updateLimitOffsets(); + break; + } + } } } } private void bufferChangelogRecords(final ChangelogMetadata changelogMetadata, final List> records) { // update the buffered records and limit index with the fetched records - final long limitOffset = Math.min( - changelogMetadata.restoreEndOffset == null ? Long.MAX_VALUE : changelogMetadata.restoreEndOffset, - changelogMetadata.restoreLimitOffset == null ? Long.MAX_VALUE : changelogMetadata.restoreLimitOffset - ); - for (final ConsumerRecord record : records) { // filter polled records for null-keys and also possibly update buffer limit index if (record.key() == null) { @@ -450,7 +453,7 @@ private void bufferChangelogRecords(final ChangelogMetadata changelogMetadata, f } else { changelogMetadata.bufferedRecords.add(record); final long offset = record.offset(); - if (offset < limitOffset) + if (changelogMetadata.restoreEndOffset == null || offset < changelogMetadata.restoreEndOffset) changelogMetadata.bufferedLimitIndex = changelogMetadata.bufferedRecords.size(); } } @@ -517,9 +520,10 @@ private Map committedOffsetForChangelogs(final Set committedOffsets; try { // those do not have a committed offset would default to 0 - return mainConsumer.committed(partitions).entrySet().stream() + committedOffsets = mainConsumer.committed(partitions).entrySet().stream() .collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue() == null ? 0L : e.getValue().offset())); } catch (final TimeoutException e) { // if it timed out we just retry next time. @@ -527,6 +531,10 @@ private Map committedOffsetForChangelogs(final Set endOffsetForChangelogs(final Set partitions) { @@ -544,7 +552,10 @@ private Map endOffsetForChangelogs(final Set committedOffsets) { @@ -568,18 +577,18 @@ private void updateLimitOffsetsForStandbyChangelogs(final Map newLimit) { throw new IllegalStateException("Offset limit should monotonically increase, but was reduced for partition " + partition + ". New limit: " + newLimit + ". Previous limit: " + previousLimit); } - metadata.restoreLimitOffset = newLimit; + metadata.restoreEndOffset = newLimit; // update the limit index for buffered records while (metadata.bufferedLimitIndex < metadata.bufferedRecords.size() && - metadata.bufferedRecords.get(metadata.bufferedLimitIndex).offset() < metadata.restoreLimitOffset) + metadata.bufferedRecords.get(metadata.bufferedLimitIndex).offset() < metadata.restoreEndOffset) metadata.bufferedLimitIndex++; } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index a1b058e97deab..4600ca1debd5d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -545,7 +545,7 @@ public static StreamThread create(final InternalTopologyBuilder builder, final Map restoreConsumerConfigs = config.getRestoreConsumerConfigs(getRestoreConsumerClientId(threadId)); final Consumer restoreConsumer = clientSupplier.getRestoreConsumer(restoreConsumerConfigs); - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, restoreConsumer, userStateRestoreListener); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, restoreConsumer, userStateRestoreListener); final ThreadCache cache = new ThreadCache(logContext, cacheSizeBytes, streamsMetrics); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/MockChangelogReader.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/MockChangelogReader.java index bcba49a2037ef..93ebeda0e0232 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/MockChangelogReader.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/MockChangelogReader.java @@ -52,11 +52,6 @@ public void transitToUpdateStandby() { // do nothing } - @Override - public void updateLimitOffsets() { - // do nothing - } - @Override public Set completedChangelogs() { // assuming all restoring partitions are completed diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java index 5a577aa7365fa..e9784eaea10b0 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StoreChangelogReaderTest.java @@ -25,6 +25,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.StateStore; @@ -46,6 +47,7 @@ import java.util.Collection; import java.util.Collections; import java.util.Map; +import java.util.Properties; import java.util.Set; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; @@ -102,6 +104,7 @@ public static Object[] data() { private final TopicPartition tp1 = new TopicPartition("one", 0); private final TopicPartition tp2 = new TopicPartition("two", 0); private final StreamsConfig config = new StreamsConfig(StreamsTestUtils.getStreamsConfig("test-reader")); + private final MockTime time = new MockTime(); private final MockStateRestoreListener callback = new MockStateRestoreListener(); private final KafkaException kaboom = new KafkaException("KABOOM!"); private final MockStateRestoreListener exceptionCallback = new MockStateRestoreListener() { @@ -127,7 +130,7 @@ public void onRestoreEnd(final TopicPartition tp, final String store, final long }; private final MockConsumer consumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); - private final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + private final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); @Before public void setUp() { @@ -156,7 +159,6 @@ public void shouldNotRegisterSameStoreMultipleTimes() { assertEquals(StoreChangelogReader.ChangelogState.REGISTERED, changelogReader.changelogMetadata(tp).state()); assertNull(changelogReader.changelogMetadata(tp).endOffset()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(0L, changelogReader.changelogMetadata(tp).totalRestored()); assertThrows(IllegalStateException.class, () -> changelogReader.register(tp, stateManager)); @@ -182,7 +184,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, stateManager); changelogReader.restore(); @@ -192,7 +194,6 @@ public Map endOffsets(final Collection par assertEquals(type == ACTIVE ? 10L : null, changelogReader.changelogMetadata(tp).endOffset()); assertEquals(0L, changelogReader.changelogMetadata(tp).totalRestored()); assertEquals(type == ACTIVE ? Collections.singleton(tp) : Collections.emptySet(), changelogReader.completedChangelogs()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(10L, consumer.position(tp)); assertEquals(Collections.singleton(tp), consumer.paused()); @@ -216,7 +217,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, stateManager); @@ -229,7 +230,6 @@ public Map endOffsets(final Collection par assertEquals(StoreChangelogReader.ChangelogState.RESTORING, changelogReader.changelogMetadata(tp).state()); assertEquals(0L, changelogReader.changelogMetadata(tp).totalRestored()); assertTrue(changelogReader.completedChangelogs().isEmpty()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(6L, consumer.position(tp)); assertEquals(Collections.emptySet(), consumer.paused()); @@ -288,7 +288,7 @@ public Map endOffsets(final Collection par }; consumer.updateBeginningOffsets(Collections.singletonMap(tp, 5L)); - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, stateManager); @@ -300,7 +300,6 @@ public Map endOffsets(final Collection par assertEquals(StoreChangelogReader.ChangelogState.RESTORING, changelogReader.changelogMetadata(tp).state()); assertEquals(0L, changelogReader.changelogMetadata(tp).totalRestored()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(5L, consumer.position(tp)); assertEquals(Collections.emptySet(), consumer.paused()); @@ -363,7 +362,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, activeStateManager); changelogReader.restore(); @@ -372,7 +371,6 @@ public Map endOffsets(final Collection par assertEquals(0L, (long) changelogReader.changelogMetadata(tp).endOffset()); assertEquals(0L, changelogReader.changelogMetadata(tp).totalRestored()); assertEquals(Collections.singleton(tp), changelogReader.completedChangelogs()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(6L, consumer.position(tp)); assertEquals(Collections.singleton(tp), consumer.paused()); assertEquals(tp, callback.restoreTopicPartition); @@ -403,7 +401,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, activeStateManager); changelogReader.restore(); @@ -411,14 +409,12 @@ public Map endOffsets(final Collection par assertEquals(StoreChangelogReader.ChangelogState.RESTORING, changelogReader.changelogMetadata(tp).state()); assertTrue(changelogReader.completedChangelogs().isEmpty()); assertEquals(10L, (long) changelogReader.changelogMetadata(tp).endOffset()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); clearException.set(true); changelogReader.restore(); assertEquals(StoreChangelogReader.ChangelogState.COMPLETED, changelogReader.changelogMetadata(tp).state()); assertEquals(10L, (long) changelogReader.changelogMetadata(tp).endOffset()); - assertNull(changelogReader.changelogMetadata(tp).limitOffset()); assertEquals(Collections.singleton(tp), changelogReader.completedChangelogs()); assertEquals(10L, consumer.position(tp)); } @@ -440,7 +436,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, activeStateManager); @@ -471,21 +467,19 @@ public Map committed(final Set endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.register(tp, activeStateManager); @@ -533,7 +527,7 @@ public Map endOffsets(final Collection par } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); changelogReader.setMainConsumer(consumer); changelogReader.register(tp, stateManager); @@ -541,15 +535,17 @@ public Map endOffsets(final Collection par assertEquals(type == ACTIVE ? StoreChangelogReader.ChangelogState.REGISTERED : StoreChangelogReader.ChangelogState.RESTORING, changelogReader.changelogMetadata(tp).state()); - assertNull(changelogReader.changelogMetadata(tp).endOffset()); - assertEquals(type == ACTIVE ? null : 0L, changelogReader.changelogMetadata(tp).limitOffset()); + if (type == ACTIVE) { + assertNull(changelogReader.changelogMetadata(tp).endOffset()); + } else { + assertEquals(0L, (long) changelogReader.changelogMetadata(tp).endOffset()); + } assertTrue(functionCalled.get()); changelogReader.restore(); assertEquals(StoreChangelogReader.ChangelogState.RESTORING, changelogReader.changelogMetadata(tp).state()); - assertEquals(type == ACTIVE ? 10L : null, changelogReader.changelogMetadata(tp).endOffset()); - assertEquals(type == ACTIVE ? null : 0L, changelogReader.changelogMetadata(tp).limitOffset()); + assertEquals(type == ACTIVE ? 10L : 0L, (long) changelogReader.changelogMetadata(tp).endOffset()); assertEquals(6L, consumer.position(tp)); } @@ -570,7 +566,7 @@ public Map committed(final Set(topicName, 0, 6L, "key".getBytes(), "value".getBytes())); @@ -645,7 +640,11 @@ public Map committed(final Set committed(final Set(topicName, 0, 5L, "key".getBytes(), "value".getBytes())); @@ -671,8 +669,7 @@ public Map committed(final Set committed(final Set committed(final Set(topicName, 0, 15L, "key".getBytes(), "value".getBytes())); changelogReader.restore(); - assertEquals(15L, (long) changelogReader.changelogMetadata(tp).limitOffset()); + assertEquals(15L, (long) changelogReader.changelogMetadata(tp).endOffset()); assertEquals(9L, changelogReader.changelogMetadata(tp).totalRestored()); assertEquals(1, changelogReader.changelogMetadata(tp).bufferedRecords().size()); assertEquals(0, changelogReader.changelogMetadata(tp).bufferedLimitIndex()); @@ -776,7 +793,7 @@ public Map endOffsets(final Collection par return partitions.stream().collect(Collectors.toMap(Function.identity(), partition -> 10L)); } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, callback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, callback); assertEquals(ACTIVE_RESTORING, changelogReader.state()); changelogReader.register(tp, activeStateManager); @@ -843,7 +860,7 @@ public Map endOffsets(final Collection par return partitions.stream().collect(Collectors.toMap(Function.identity(), partition -> 10L)); } }; - final StoreChangelogReader changelogReader = new StoreChangelogReader(config, logContext, consumer, exceptionCallback); + final StoreChangelogReader changelogReader = new StoreChangelogReader(time, config, logContext, consumer, exceptionCallback); changelogReader.register(tp, activeStateManager); diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java index 5c6da5c4b82e9..e16b778be7246 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java @@ -318,6 +318,7 @@ private StreamTask createStreamsTask(final StreamsConfig streamsConfig, stateDirectory, topology.storeToChangelogTopic(), new StoreChangelogReader( + new MockTime(), streamsConfig, logContext, clientSupplier.restoreConsumer, diff --git a/streams/src/test/java/org/apache/kafka/test/StreamsTestUtils.java b/streams/src/test/java/org/apache/kafka/test/StreamsTestUtils.java index 588c0172c5982..50e3de74f9a09 100644 --- a/streams/src/test/java/org/apache/kafka/test/StreamsTestUtils.java +++ b/streams/src/test/java/org/apache/kafka/test/StreamsTestUtils.java @@ -74,12 +74,16 @@ public static Properties getStreamsConfig(final Serde keyDeserializer, } public static Properties getStreamsConfig(final String applicationId) { + return getStreamsConfig(applicationId, new Properties()); + } + + public static Properties getStreamsConfig(final String applicationId, final Properties additional) { return getStreamsConfig( applicationId, "localhost:9091", Serdes.ByteArraySerde.class.getName(), Serdes.ByteArraySerde.class.getName(), - new Properties()); + additional); } public static Properties getStreamsConfig() { diff --git a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java index 78f03a6f6dd52..c7ae313698ff5 100644 --- a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java +++ b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java @@ -390,6 +390,7 @@ public List partitionsFor(final String topic) { stateDirectory, processorTopology.storeToChangelogTopic(), new StoreChangelogReader( + mockWallClockTime, streamsConfig, logContext, createRestoreConsumer(processorTopology.storeToChangelogTopic()), From 33af581d8cb942c5123e4dfc019514aa99129ba8 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 6 Feb 2020 13:09:55 -0800 Subject: [PATCH 02/17] consumer.position invalid offset --- .../streams/processor/internals/StoreChangelogReader.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index dac6115c2a399..92a92ed9d649c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -258,6 +258,8 @@ private boolean hasRestoredToEnd(final ChangelogMetadata metadata) { // if we cannot get the position of the consumer within timeout, just return false return false; } catch (final KafkaException e) { + // this also includes InvalidOffsetException, which should not happen under normal + // execution, hence it is also okay to wrap it as fatal StreamsException throw new StreamsException("Restore consumer get unexpected error trying to get the position " + " of " + partition, e); } @@ -755,6 +757,8 @@ private void prepareChangelogs(final Set newPartitionsToResto } catch (final TimeoutException e) { // if we cannot find the starting position at the beginning, just use the default 0L } catch (final KafkaException e) { + // this also includes InvalidOffsetException, which should not happen under normal + // execution, hence it is also okay to wrap it as fatal StreamsException throw new StreamsException("Restore consumer get unexpected error trying to get the position " + " of " + partition, e); } From 4419900a7f26906a60c0658f4290cd7a2e2d07c6 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 6 Feb 2020 15:14:09 -0800 Subject: [PATCH 03/17] add task corrupted logic --- .../errors/TaskCorruptedException.java | 44 +++++++++++ .../streams/errors/TaskMigratedException.java | 41 +--------- .../processor/internals/AbstractTask.java | 10 +++ .../internals/RecordCollectorImpl.java | 10 +-- .../internals/StoreChangelogReader.java | 38 +++++++--- .../processor/internals/StreamTask.java | 15 ++-- .../processor/internals/StreamThread.java | 39 ++++++---- .../streams/processor/internals/Task.java | 55 ++++++++------ .../processor/internals/TaskManager.java | 76 ++++++++++--------- .../processor/internals/TaskManagerTest.java | 7 +- 10 files changed, 193 insertions(+), 142 deletions(-) create mode 100644 streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java diff --git a/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java b/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java new file mode 100644 index 0000000000000..e63c9c0e86e85 --- /dev/null +++ b/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java @@ -0,0 +1,44 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.streams.errors; + +import org.apache.kafka.streams.processor.TaskId; + +import java.util.Set; + +/** + * Indicates a specific task is corrupted and need to be re-initialized. It can be thrown when + * + * 1) Under EOS, if the checkpoint file does not contain offsets for corresponding store's changelogs, meaning + * previously it was not close cleanly; + * 2) Out-of-range exception thrown during restoration, meaning that the changelog has been modified and we re-bootstrap + * the store. + */ +public class TaskCorruptedException extends StreamsException { + + private final Set taskIds; + + public TaskCorruptedException(final Set taskIds) { + super("Tasks " + taskIds + " are corrupted and hence needs to be re-initialized"); + + this.taskIds = taskIds; + } + + public Set corruptedTaskIds() { + return taskIds; + } +} diff --git a/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java b/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java index a3307f55a2a54..d60321e17b064 100644 --- a/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java +++ b/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java @@ -17,9 +17,6 @@ package org.apache.kafka.streams.errors; -import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.streams.processor.TaskId; - /** * Indicates that one or more tasks got migrated to another thread. * @@ -32,43 +29,7 @@ public class TaskMigratedException extends StreamsException { private final static long serialVersionUID = 1L; - private final TaskId taskId; - - public TaskMigratedException(final TaskId taskId, - final TopicPartition topicPartition, - final long endOffset, - final long pos) { - this(taskId, String.format("Log end offset of %s should not change while restoring: old end offset %d, current offset %d", - topicPartition, - endOffset, - pos), null); - } - - public TaskMigratedException(final TaskId taskId) { - this(taskId, String.format("Task %s is unexpectedly closed during processing", taskId), null); - } - - public TaskMigratedException(final TaskId taskId, - final Throwable throwable) { - this(taskId, String.format("Client request for task %s has been fenced due to a rebalance", taskId), throwable); - } - - public TaskMigratedException(final TaskId taskId, - final String message, - final Throwable throwable) { - super(message, throwable); - this.taskId = taskId; - } - public TaskMigratedException(final String message, final Throwable throwable) { - this(null, message + " It means all tasks belonging to this thread have been migrated", throwable); - } - - public TaskId migratedTaskId() { - return taskId; - } - - public TaskMigratedException() { - this(null, "A task has been migrated unexpectedly", null); + super(message + "; It means all tasks belonging to this thread have been migrated", throwable); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java index ca687316d5e2e..dcb4b5f71f060 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java @@ -22,6 +22,7 @@ import java.util.Set; +import static org.apache.kafka.streams.processor.internals.Task.State.CLOSED; import static org.apache.kafka.streams.processor.internals.Task.State.CREATED; public abstract class AbstractTask implements Task { @@ -70,6 +71,15 @@ public final Task.State state() { return state; } + @Override + public void revive() { + if (state == CLOSED) { + transitionTo(CREATED); + } else { + throw new IllegalStateException("Illegal state " + state() + " while committing standby task " + id); + } + } + final void transitionTo(final Task.State newState) { final State oldState = state(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordCollectorImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordCollectorImpl.java index a6399858fa74a..4cea3224d1064 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordCollectorImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/RecordCollectorImpl.java @@ -124,7 +124,7 @@ private void maybeBeginTxn() { try { producer.beginTransaction(); } catch (final ProducerFencedException error) { - throw new TaskMigratedException(taskId, "Producer get fenced trying to begin a new transaction", error); + throw new TaskMigratedException("Producer get fenced trying to begin a new transaction", error); } catch (final KafkaException error) { throw new StreamsException("Producer encounter unexpected error trying to begin a new transaction", error); } @@ -161,7 +161,7 @@ public void commit(final Map offsets) { producer.commitTransaction(); transactionInFlight = false; } catch (final ProducerFencedException error) { - throw new TaskMigratedException(taskId, "Producer get fenced trying to commit a transaction", error); + throw new TaskMigratedException("Producer get fenced trying to commit a transaction", error); } catch (final TimeoutException error) { // TODO K9113: currently handle timeout exception as a fatal error, should discuss whether we want to handle it throw new StreamsException("Timed out while committing transaction via producer for task " + taskId, error); @@ -173,7 +173,7 @@ public void commit(final Map offsets) { try { consumer.commitSync(offsets); } catch (final CommitFailedException error) { - throw new TaskMigratedException(taskId, "Consumer committing offsets failed, " + + throw new TaskMigratedException("Consumer committing offsets failed, " + "indicating the corresponding thread is no longer part of the group.", error); } catch (final TimeoutException error) { // TODO K9113: currently handle timeout exception as a fatal error @@ -208,7 +208,7 @@ private void recordSendError(final String topic, final Exception exception, fina } else if (exception instanceof ProducerFencedException || exception instanceof OutOfOrderSequenceException) { errorMessage += "\nWritten offsets would not be recorded and no more records would be sent since the producer is fenced, " + "indicating the task may be migrated out."; - sendException = new TaskMigratedException(taskId, errorMessage, exception); + sendException = new TaskMigratedException(errorMessage, exception); } else { if (exception instanceof RetriableException) { errorMessage += "\nThe broker is either slow or in bad state (like not having enough replicas) in responding the request, " + @@ -299,7 +299,7 @@ public void send(final String topic, if (isRecoverable(uncaughtException)) { // producer.send() call may throw a KafkaException which wraps a FencedException, // in this case we should throw its wrapped inner cause so that it can be captured and re-wrapped as TaskMigrationException - throw new TaskMigratedException(taskId, "Producer cannot send records anymore since it got fenced", uncaughtException.getCause()); + throw new TaskMigratedException("Producer cannot send records anymore since it got fenced", uncaughtException.getCause()); } else { final String errorMessage = String.format(SEND_EXCEPTION_MESSAGE, topic, taskId, uncaughtException.toString()); throw new StreamsException(errorMessage, uncaughtException); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index 92a92ed9d649c..f3df7c4603b8c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -19,6 +19,7 @@ import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.InvalidOffsetException; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.errors.FencedInstanceIdException; import org.apache.kafka.common.errors.TimeoutException; @@ -27,7 +28,9 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.errors.TaskCorruptedException; import org.apache.kafka.streams.errors.TaskMigratedException; +import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.ProcessorStateManager.StateStoreMetadata; import org.apache.kafka.streams.processor.StateRestoreListener; import org.slf4j.Logger; @@ -52,7 +55,6 @@ * The reader also maintains the source of truth for restoration state: only active tasks restoring changelog could * be completed, while standby tasks updating changelog would always be in restoring state after being initialized. */ -// TODO K9113: we need to consider how to handle InvalidOffsetException for consumer#poll / position public class StoreChangelogReader implements ChangelogReader { enum ChangelogState { @@ -411,6 +413,15 @@ public void restore() { } catch (final FencedInstanceIdException e) { // when the consumer gets fenced, all its tasks should be migrated throw new TaskMigratedException("Restore consumer get fenced by instance-id polling records.", e); + } catch (final InvalidOffsetException e) { + log.warn("Encountered {} fetching records from restore consumer for partitions {}, " + + "marking the corresponding tasks as corrupted.", e.toString(), e.partitions()); + + final Set taskIds = new HashSet<>(); + for (final TopicPartition partition : e.partitions()) { + taskIds.add(changelogs.get(partition).stateManager.taskId()); + } + throw new TaskCorruptedException(taskIds); } catch (final KafkaException e) { throw new StreamsException("Restore consumer get unexpected error polling records.", e); } @@ -428,18 +439,21 @@ public void restore() { restoreChangelog(changelogs.get(partition)); } + maybeUpdateLimitOffsetsForStandbyChangelogs(); + } + } - // for standby changelogs, if the interval has elapsed and there are buffered records not applicable, - // we can try to update the limit offset next time. - if (updateOffsetIntervalMs < time.milliseconds() - lastUpdateOffsetTime) { - final Set standbyChangelogs = changelogs.values().stream() - .filter(metadata -> metadata.stateManager.taskType() == Task.TaskType.STANDBY) - .collect(Collectors.toSet()); - for (final ChangelogMetadata metadata : standbyChangelogs) { - if (!metadata.bufferedRecords().isEmpty()) { - updateLimitOffsets(); - break; - } + private void maybeUpdateLimitOffsetsForStandbyChangelogs() { + // for standby changelogs, if the interval has elapsed and there are buffered records not applicable, + // we can try to update the limit offset next time. + if (updateOffsetIntervalMs < time.milliseconds() - lastUpdateOffsetTime) { + final Set standbyChangelogs = changelogs.values().stream() + .filter(metadata -> metadata.stateManager.taskType() == Task.TaskType.STANDBY) + .collect(Collectors.toSet()); + for (final ChangelogMetadata metadata : standbyChangelogs) { + if (!metadata.bufferedRecords().isEmpty()) { + updateLimitOffsets(); + break; } } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 8cd6e5b4ead75..54727a3bcf50e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -69,10 +69,9 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator, // visible for testing static final byte LATEST_MAGIC_BYTE = 1; + private final Time time; private final Logger log; private final String logPrefix; - private final Time time; - private final String threadId; private final Consumer consumer; // we want to abstract eos logic out of StreamTask, however @@ -82,7 +81,6 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator, private final long maxTaskIdleMs; private final int maxBufferedSize; - private final StreamsMetricsImpl streamsMetrics; private final PartitionGroup partitionGroup; private final RecordCollector recordCollector; private final PartitionGroup.RecordInfo recordInfo; @@ -121,11 +119,10 @@ public StreamTask(final TaskId id, log = logContext.logger(getClass()); this.time = time; - this.streamsMetrics = streamsMetrics; this.recordCollector = recordCollector; eosDisabled = !StreamsConfig.EXACTLY_ONCE.equals(config.getString(StreamsConfig.PROCESSING_GUARANTEE_CONFIG)); - threadId = Thread.currentThread().getName(); + final String threadId = Thread.currentThread().getName(); closeTaskSensor = ThreadMetrics.closeTaskSensor(threadId, streamsMetrics); final String taskId = id.toString(); if (streamsMetrics.version() == Version.FROM_0100_TO_24) { @@ -435,7 +432,6 @@ private void close(final boolean clean) { partitionGroup.close(); closeTaskSensor.record(); - streamsMetrics.removeAllTaskLevelSensors(threadId, id.toString()); transitionTo(State.CLOSED); } @@ -692,9 +688,10 @@ private void closeRecordCollector(final boolean clean) { @Override public void addRecords(final TopicPartition partition, final Iterable> records) { if (state() == State.CLOSED || state() == State.CLOSING) { - log.info("Stream task {} is already closed, probably because it got unexpectedly migrated to another thread already. " + - "Notifying the thread to trigger a new rebalance immediately.", id()); - throw new TaskMigratedException(id()); + // a task is only closing / closed when 1) task manager is closing, 2) a rebalance is undergoing; + // in either case we can just log it and move on without notifying the thread since the consumer + // would soon be updated to not return any records for this task anymore. + log.info("Stream task {} is already in {} state, skip adding records to it.", id(), state()); } final int newQueueSize = partitionGroup.addRawRecords(partition, records); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 4600ca1debd5d..5cf81c56d71db 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -35,6 +35,7 @@ import org.apache.kafka.streams.KafkaClientSupplier; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.errors.TaskCorruptedException; import org.apache.kafka.streams.errors.TaskMigratedException; import org.apache.kafka.streams.processor.StateRestoreListener; import org.apache.kafka.streams.processor.TaskId; @@ -262,7 +263,6 @@ static abstract class AbstractTaskCreator { final Time time; final Logger log; - AbstractTaskCreator(final InternalTopologyBuilder builder, final StreamsConfig config, final StreamsMetricsImpl streamsMetrics, @@ -573,6 +573,7 @@ public static StreamThread create(final InternalTopologyBuilder builder, changelogReader, processId, logPrefix, + streamsMetrics, activeTaskCreator, standbyTaskCreator, builder, @@ -718,14 +719,19 @@ public void run() { try { runLoop(); cleanRun = true; + } catch (final StreamsException e) { + // do not need to rethrow + log.error("Encountered the following Streams exception during processing, " + + "this indicates an unrecoverable error and the thread is going to shutdown: ", e); } catch (final KafkaException e) { log.error("Encountered the following unexpected Kafka exception during processing, " + - "this usually indicate Streams internal errors:", e); + "this indicates an internal error and thread is going to shutdown:", e); throw e; } catch (final Exception e) { // we have caught all Kafka related exceptions, and other runtime exceptions // should be due to user application errors - log.error("Encountered the following error during processing:", e); + log.error("Encountered the following unexpected exception during processing " + + "and the thread is going to shutdown:", e); throw e; } finally { completeShutdown(cleanRun); @@ -745,14 +751,23 @@ private void runLoop() { try { runOnce(); if (assignmentErrorCode.get() == AssignorError.VERSION_PROBING.code()) { - log.info("Version probing detected. Triggering new rebalance."); + log.info("Version probing detected. Rejoin the consumer group to trigger a new rebalance."); assignmentErrorCode.set(AssignorError.NONE.code()); enforceRebalance(); } - } catch (final TaskMigratedException ignoreAndRejoinGroup) { - log.warn("Detected task {} that got migrated to another thread. " + - "This implies that this thread missed a rebalance and dropped out of the consumer group. " + - "Will try to rejoin the consumer group.", ignoreAndRejoinGroup.migratedTaskId()); + } catch (final TaskCorruptedException e) { + log.warn("Detected the states of tasks {} are corrupted. " + + "Will close the task as dirty and re-create and bootstrap from scratch."); + + for (final TaskId taskId : e.corruptedTaskIds()) { + final Task task = taskManager.tasks().get(taskId); + task.closeDirty(); + task.revive(); + } + } catch (final TaskMigratedException e) { + log.warn("Detected that the thread is being fenced. " + + "This implies that this thread missed a rebalance and dropped out of the consumer group. " + + "Will migrate out all assigned tasks and rejoin the consumer group."); enforceRebalance(); } @@ -964,16 +979,12 @@ private void addToResetList(final TopicPartition partition, final Set records) { - for (final TopicPartition partition : records.partitions()) { final Task task = taskManager.taskForInputPartition(partition); if (task == null) { - log.error( - "Unable to locate active task for received-record partition {}. Current tasks: {}", - partition, - taskManager.toString(">") - ); + log.error("Unable to locate active task for received-record partition {}. Current tasks: {}", + partition, taskManager.toString(">")); throw new NullPointerException("Task was unexpectedly missing for partition " + partition); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java index e91b885f93075..57fa18eb8dfb2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java @@ -39,30 +39,30 @@ public interface Task { * *
      *                 +-------------+
-     *          +<---- | Created (0) |
-     *          |      +-----+-------+
-     *          |            |
-     *          |            v
-     *          |      +-----+-------+
-     *          +<---- | Restoring(1)|<---------------+
-     *          |      +-----+-------+                |
-     *          |            |                        |
-     *          |            +--------------------+   |
-     *          |            |                    |   |
-     *          |            v                    v   |
-     *          |      +-----+-------+       +----+---+----+
-     *          |      | Running (2) | ----> | Suspended(3)|   * //TODO Suspended(3) could be removed after we've stable on KIP-429
-     *          |      +-----+-------+       +------+------+
-     *          |            |                      |
-     *          |            |                      |
-     *          |            v                      |
-     *          |      +-----+-------+              |
-     *          +----> | Closing (4) | <------------+
-     *                 +-----+-------+
-     *                       |
-     *                       v
-     *                 +-----+-------+
-     *                 | Closed (5)  |
+     *          +<---- | Created (0) | <----------------------+
+     *          |      +-----+-------+                        |
+     *          |            |                                |
+     *          |            v                                |
+     *          |      +-----+-------+                        |
+     *          +<---- | Restoring(1)|<---------------+       |
+     *          |      +-----+-------+                |       |
+     *          |            |                        |       |
+     *          |            +--------------------+   |       |
+     *          |            |                    |   |       |
+     *          |            v                    v   |       |
+     *          |      +-----+-------+       +----+---+----+  |
+     *          |      | Running (2) | ----> | Suspended(3)|  |    * //TODO Suspended(3) could be removed after we've stable on KIP-429
+     *          |      +-----+-------+       +------+------+  |
+     *          |            |                      |         |
+     *          |            |                      |         |
+     *          |            v                      |         |
+     *          |      +-----+-------+              |         |
+     *          +----> | Closing (4) | <------------+         |
+     *                 +-----+-------+                        |
+     *                       |                                |
+     *                       v                                |
+     *                 +-----+-------+                        |
+     *                 | Closed (5)  | -----------------------+
      *                 +-------------+
      * 
*/ @@ -72,7 +72,7 @@ enum State { RUNNING(3, 4), // 2 SUSPENDED(1, 4), // 3 CLOSING(4, 5), // 4, we allow CLOSING to transit to itself to make close idempotent - CLOSED; // 5 + CLOSED(0); // 5, we allow CLOSED to transit to CREATED to handle corrupted tasks private final Set validTransitions = new HashSet<>(); @@ -154,6 +154,11 @@ enum TaskType { */ void closeDirty(); + /** + * Revive a closed task to a created one; should never throw an exception + */ + void revive(); + StateStore getStore(final String name); Set inputPartitions(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index 4b151ce208d62..e1cefcabc2b00 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -28,6 +28,7 @@ import org.apache.kafka.streams.errors.TaskIdFormatException; import org.apache.kafka.streams.errors.TaskMigratedException; import org.apache.kafka.streams.processor.TaskId; +import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.slf4j.Logger; import java.io.File; @@ -55,40 +56,42 @@ public class TaskManager { // by QueryableState private final Logger log; private final UUID processId; - private final ChangelogReader changelogReader; private final String logPrefix; - private final StreamThread.AbstractTaskCreator taskCreator; + private final InternalTopologyBuilder builder; + private final ChangelogReader changelogReader; + private final StreamsMetricsImpl streamsMetrics; + private final StreamThread.AbstractTaskCreator activeTaskCreator; private final StreamThread.AbstractTaskCreator standbyTaskCreator; - private final Admin adminClient; - private DeleteRecordsResult deleteRecordsResult; - private boolean rebalanceInProgress = false; // if we are in the middle of a rebalance, it is not safe to commit - private final Map tasks = new TreeMap<>(); // materializing this relationship because the lookup is on the hot path private final Map partitionToTask = new HashMap<>(); + private final Admin adminClient; private Consumer consumer; - private final InternalTopologyBuilder builder; + private DeleteRecordsResult deleteRecordsResult; + + private boolean rebalanceInProgress = false; // if we are in the middle of a rebalance, it is not safe to commit TaskManager(final ChangelogReader changelogReader, final UUID processId, final String logPrefix, - final StreamThread.AbstractTaskCreator taskCreator, + final StreamsMetricsImpl streamsMetrics, + final StreamThread.AbstractTaskCreator activeTaskCreator, final StreamThread.AbstractTaskCreator standbyTaskCreator, final InternalTopologyBuilder builder, final Admin adminClient) { - this.changelogReader = changelogReader; + this.builder = builder; this.processId = processId; this.logPrefix = logPrefix; - this.taskCreator = taskCreator; + this.adminClient = adminClient; + this.streamsMetrics = streamsMetrics; + this.changelogReader = changelogReader; + this.activeTaskCreator = activeTaskCreator; this.standbyTaskCreator = standbyTaskCreator; - this.builder = builder; - final LogContext logContext = new LogContext(logPrefix); - - log = logContext.logger(getClass()); - this.adminClient = adminClient; + final LogContext logContext = new LogContext(logPrefix); + this.log = logContext.logger(getClass()); } void setConsumer(final Consumer consumer) { @@ -147,10 +150,8 @@ public void handleAssignment(final Map> activeTasks, task.resume(); standbyTasksToCreate.remove(task.id()); } else /* we previously owned this task, and we don't have it anymore, or it has changed active/standby state */ { - final Set inputPartitions = task.inputPartitions(); try { task.closeClean(); - changelogReader.remove(task.changelogPartitions()); } catch (final RuntimeException e) { log.error("Failed to close task {} cleanly. Attempting to close remaining tasks before re-throwing.", task.id()); taskCloseExceptions.put(task.id(), e); @@ -158,9 +159,8 @@ public void handleAssignment(final Map> activeTasks, // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. task.closeDirty(); } - for (final TopicPartition inputPartition : inputPartitions) { - partitionToTask.remove(inputPartition); - } + cleanupTask(task); + iterator.remove(); } } @@ -175,7 +175,7 @@ public void handleAssignment(final Map> activeTasks, } if (!activeTasksToCreate.isEmpty()) { - taskCreator.createTasks(consumer, activeTasksToCreate).forEach(this::addNewTask); + activeTaskCreator.createTasks(consumer, activeTasksToCreate).forEach(this::addNewTask); } if (!standbyTasksToCreate.isEmpty()) { @@ -285,18 +285,13 @@ void handleLostAll() { final Iterator iterator = tasks.values().iterator(); while (iterator.hasNext()) { final Task task = iterator.next(); - final Set inputPartitions = task.inputPartitions(); // Even though we've apparently dropped out of the group, we can continue safely to maintain our // standby tasks while we rejoin. if (task.isActive()) { task.closeDirty(); - changelogReader.remove(task.changelogPartitions()); - } - - for (final TopicPartition inputPartition : inputPartitions) { - partitionToTask.remove(inputPartition); + cleanupTask(task); + iterator.remove(); } - iterator.remove(); } } @@ -312,7 +307,7 @@ public Set tasksOnLocalStorage() { final Set locallyStoredTasks = new HashSet<>(); - final File[] stateDirs = taskCreator.stateDirectory().listTaskDirectories(); + final File[] stateDirs = activeTaskCreator.stateDirectory().listTaskDirectories(); if (stateDirs != null) { for (final File dir : stateDirs) { try { @@ -331,12 +326,25 @@ public Set tasksOnLocalStorage() { return locallyStoredTasks; } + private void cleanupTask(final Task task) { + // 1. remove the changelog partitions from changelog reader; + // 2. remove the input partitions from the materialized map; + // 3. remove the task metrics from the metrics registry + changelogReader.remove(task.changelogPartitions()); + + for (final TopicPartition inputPartition : task.inputPartitions()) { + partitionToTask.remove(inputPartition); + } + + final String threadId = Thread.currentThread().getName(); + streamsMetrics.removeAllTaskLevelSensors(threadId, task.id().toString()); + } + void shutdown(final boolean clean) { final AtomicReference firstException = new AtomicReference<>(null); final Iterator iterator = tasks.values().iterator(); while (iterator.hasNext()) { final Task task = iterator.next(); - final Set inputPartitions = task.inputPartitions(); if (clean) { try { task.closeClean(); @@ -350,16 +358,12 @@ void shutdown(final boolean clean) { } else { task.closeDirty(); } - changelogReader.remove(task.changelogPartitions()); - - for (final TopicPartition inputPartition : inputPartitions) { - partitionToTask.remove(inputPartition); - } + cleanupTask(task); iterator.remove(); } - taskCreator.close(); + activeTaskCreator.close(); final RuntimeException fatalException = firstException.get(); if (fatalException != null) { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java index dbc237a56ec34..a8f262ca71b9b 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/TaskManagerTest.java @@ -25,9 +25,12 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.internals.KafkaFutureImpl; +import org.apache.kafka.common.metrics.Metrics; +import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; +import org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl; import org.easymock.EasyMock; import org.easymock.EasyMockRunner; import org.easymock.Mock; @@ -115,9 +118,11 @@ public class TaskManagerTest { @Before public void setUp() { + final StreamsMetricsImpl streamsMetrics = new StreamsMetricsImpl(new Metrics(), "clientId", StreamsConfig.METRICS_LATEST); taskManager = new TaskManager(changeLogReader, UUID.randomUUID(), - "", + "taskManagerTest", + streamsMetrics, activeTaskCreator, standbyTaskCreator, topologyBuilder, From 23cd292d4497e5c9d64c6c2efea69e570b72791a Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 6 Feb 2020 16:31:57 -0800 Subject: [PATCH 04/17] re-enable the unit test --- .../kafka/common/utils/FixedOrderMap.java | 6 ----- .../internals/ProcessorStateManager.java | 2 ++ .../internals/StoreChangelogReader.java | 18 +++++++++------ .../processor/internals/StreamThread.java | 14 +++++------- .../processor/internals/TaskManager.java | 22 +++++++++++++++---- .../processor/internals/StreamThreadTest.java | 6 ++--- 6 files changed, 39 insertions(+), 29 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/FixedOrderMap.java b/clients/src/main/java/org/apache/kafka/common/utils/FixedOrderMap.java index 6518878877028..175282e7e09f6 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/FixedOrderMap.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/FixedOrderMap.java @@ -50,12 +50,6 @@ public boolean remove(final Object key, final Object value) { throw new UnsupportedOperationException("Removing from registeredStores is not allowed"); } - @Deprecated - @Override - public void clear() { - throw new UnsupportedOperationException("Removing from registeredStores is not allowed"); - } - @Override public FixedOrderMap clone() { throw new UnsupportedOperationException(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java index c50cf7174ff97..ab29938030edf 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java @@ -415,6 +415,8 @@ public void close() throws ProcessorStateException { log.error("Failed to close state store {}: ", store.name(), exception); } } + + stores.clear(); } if (firstException != null) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index f3df7c4603b8c..4b96dd12529c3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -415,7 +415,7 @@ public void restore() { throw new TaskMigratedException("Restore consumer get fenced by instance-id polling records.", e); } catch (final InvalidOffsetException e) { log.warn("Encountered {} fetching records from restore consumer for partitions {}, " + - "marking the corresponding tasks as corrupted.", e.toString(), e.partitions()); + "marking the corresponding tasks as corrupted.", e.getClass().getName(), e.partitions()); final Set taskIds = new HashSet<>(); for (final TopicPartition partition : e.partitions()) { @@ -691,6 +691,8 @@ private void addChangelogsToRestoreConsumer(final Set partitions } assignment.addAll(partitions); restoreConsumer.assign(assignment); + + log.debug("Added partitions {} to the restore consumer, current assignment is {}", partitions, assignment); } private void pauseChangelogsFromRestoreConsumer(final Collection partitions) { @@ -702,18 +704,18 @@ private void pauseChangelogsFromRestoreConsumer(final Collection "does not contain some of the partitions " + partitions + " for pausing."); } restoreConsumer.pause(partitions); + + log.debug("Paused partitions {} from the restore consumer", partitions); } private void removeChangelogsFromRestoreConsumer(final Collection partitions) { final Set assignment = new HashSet<>(restoreConsumer.assignment()); - // the current assignment should contain the all partitions to remove - if (!assignment.containsAll(partitions)) { - throw new IllegalStateException("The current assignment " + assignment + " " + - "does not contain some of the partitions " + partitions + " for removing."); + if (assignment.removeAll(partitions)) { + restoreConsumer.assign(assignment); + + log.debug("Removed partitions {} from the restore consumer, current assignment is {}", partitions, assignment); } - assignment.removeAll(partitions); - restoreConsumer.assign(assignment); } private void resumeChangelogsFromRestoreConsumer(final Collection partitions) { @@ -725,6 +727,8 @@ private void resumeChangelogsFromRestoreConsumer(final Collection newPartitionsToRestore) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 5cf81c56d71db..a9b82d2cc8573 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -723,6 +723,7 @@ public void run() { // do not need to rethrow log.error("Encountered the following Streams exception during processing, " + "this indicates an unrecoverable error and the thread is going to shutdown: ", e); + throw e; } catch (final KafkaException e) { log.error("Encountered the following unexpected Kafka exception during processing, " + "this indicates an internal error and thread is going to shutdown:", e); @@ -752,18 +753,15 @@ private void runLoop() { runOnce(); if (assignmentErrorCode.get() == AssignorError.VERSION_PROBING.code()) { log.info("Version probing detected. Rejoin the consumer group to trigger a new rebalance."); + assignmentErrorCode.set(AssignorError.NONE.code()); enforceRebalance(); } } catch (final TaskCorruptedException e) { log.warn("Detected the states of tasks {} are corrupted. " + - "Will close the task as dirty and re-create and bootstrap from scratch."); + "Will close the task as dirty and re-create and bootstrap from scratch.", e.corruptedTaskIds()); - for (final TaskId taskId : e.corruptedTaskIds()) { - final Task task = taskManager.tasks().get(taskId); - task.closeDirty(); - task.revive(); - } + taskManager.handleCorruption(e.corruptedTaskIds()); } catch (final TaskMigratedException e) { log.warn("Detected that the thread is being fenced. " + "This implies that this thread missed a rebalance and dropped out of the consumer group. " + @@ -862,7 +860,7 @@ void runOnce() { * 5. If one of the above happens, half the value of N. */ int processed = 0; - long timeSinceLastPoll = 0L; + long timeSinceLastPoll; do { for (int i = 0; i < numIterations; i++) { @@ -951,7 +949,7 @@ private void resetInvalidOffsets(final InvalidOffsetException e) { if (originalReset.equals("earliest")) { addToResetList(partition, seekToBeginning, "No custom setting defined for topic '{}' using original config '{}' for offset reset", "earliest", loggedTopics); - } else if (originalReset.equals("latest")) { + } else { addToResetList(partition, seekToEnd, "No custom setting defined for topic '{}' using original config '{}' for offset reset", "latest", loggedTopics); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index e1cefcabc2b00..2e5ddf79641af 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -120,6 +120,19 @@ void handleRebalanceComplete() { rebalanceInProgress = false; } + void handleCorruption(final Set taskIds) { + for (final TaskId taskId : taskIds) { + final Task task = tasks.get(taskId); + + // this call is idempotent so even if the task is only CREATED we can still call it + changelogReader.remove(task.changelogPartitions()); + + task.closeDirty(); + + task.revive(); + } + } + /** * @throws TaskMigratedException if the task producer got fenced (EOS only) * @throws StreamsException fatal error while creating / initializing the task @@ -150,6 +163,8 @@ public void handleAssignment(final Map> activeTasks, task.resume(); standbyTasksToCreate.remove(task.id()); } else /* we previously owned this task, and we don't have it anymore, or it has changed active/standby state */ { + cleanupTask(task); + try { task.closeClean(); } catch (final RuntimeException e) { @@ -159,7 +174,6 @@ public void handleAssignment(final Map> activeTasks, // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. task.closeDirty(); } - cleanupTask(task); iterator.remove(); } @@ -288,8 +302,8 @@ void handleLostAll() { // Even though we've apparently dropped out of the group, we can continue safely to maintain our // standby tasks while we rejoin. if (task.isActive()) { - task.closeDirty(); cleanupTask(task); + task.closeDirty(); iterator.remove(); } } @@ -345,6 +359,8 @@ void shutdown(final boolean clean) { final Iterator iterator = tasks.values().iterator(); while (iterator.hasNext()) { final Task task = iterator.next(); + cleanupTask(task); + if (clean) { try { task.closeClean(); @@ -358,8 +374,6 @@ void shutdown(final boolean clean) { } else { task.closeDirty(); } - cleanupTask(task); - iterator.remove(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index 71364b3c9fb30..213d8c700b411 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -73,7 +73,6 @@ import org.easymock.EasyMock; import org.junit.Assert; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.slf4j.Logger; @@ -1366,12 +1365,11 @@ private void assertThreadMetadataHasEmptyTasksWithState(final ThreadMetadata met assertTrue(metadata.standbyTasks().isEmpty()); } - @Ignore @Test - // FIXME: should unblock this test after we added invalid offset handling public void shouldRecoverFromInvalidOffsetExceptionOnRestoreAndFinishRestore() throws Exception { internalStreamsBuilder.stream(Collections.singleton("topic"), consumed) - .groupByKey().count(Materialized.as("count")); + .groupByKey() + .count(Materialized.as("count")); internalStreamsBuilder.buildAndOptimizeTopology(); final StreamThread thread = createStreamThread("clientId", config, false); From 64dad101c5c2d8282210c0e7b431fe1f4d5f036b Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Fri, 7 Feb 2020 09:15:54 -0800 Subject: [PATCH 05/17] PR comments --- .../apache/kafka/streams/processor/internals/AbstractTask.java | 2 +- .../apache/kafka/streams/processor/internals/StreamThread.java | 3 ++- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java index dcb4b5f71f060..fc98117eb47a8 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java @@ -76,7 +76,7 @@ public void revive() { if (state == CLOSED) { transitionTo(CREATED); } else { - throw new IllegalStateException("Illegal state " + state() + " while committing standby task " + id); + throw new IllegalStateException("Illegal state " + state() + " while reviving task " + id); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index a9b82d2cc8573..7bedc40a203d3 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -765,8 +765,9 @@ private void runLoop() { } catch (final TaskMigratedException e) { log.warn("Detected that the thread is being fenced. " + "This implies that this thread missed a rebalance and dropped out of the consumer group. " + - "Will migrate out all assigned tasks and rejoin the consumer group."); + "Will close out all assigned tasks and rejoin the consumer group."); + taskManager.handleLostAll(); enforceRebalance(); } } From 134dda3187ff41addddcb7c022ff34cbea0154b3 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Fri, 7 Feb 2020 13:30:16 -0800 Subject: [PATCH 06/17] remove FixedOrderMapTest --- .../kafka/common/utils/FixedOrderMapTest.java | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/clients/src/test/java/org/apache/kafka/common/utils/FixedOrderMapTest.java b/clients/src/test/java/org/apache/kafka/common/utils/FixedOrderMapTest.java index 7d3f3f7525677..d4b41519d79d2 100644 --- a/clients/src/test/java/org/apache/kafka/common/utils/FixedOrderMapTest.java +++ b/clients/src/test/java/org/apache/kafka/common/utils/FixedOrderMapTest.java @@ -69,18 +69,4 @@ public void shouldForbidConditionalRemove() { } assertThat(map.get("a"), is(0)); } - - @SuppressWarnings("deprecation") - @Test - public void shouldForbidConditionalClear() { - final FixedOrderMap map = new FixedOrderMap<>(); - map.put("a", 0); - try { - map.clear(); - fail("expected exception"); - } catch (final RuntimeException e) { - assertThat(e, CoreMatchers.instanceOf(UnsupportedOperationException.class)); - } - assertThat(map.get("a"), is(0)); - } } From a8d81c0dbb67c1da50e23bec0f49421c6f491a49 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Sat, 8 Feb 2020 10:28:39 -0800 Subject: [PATCH 07/17] address comments --- .../streams/processor/internals/StoreChangelogReader.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index bbd1dc7cf8601..889e9332c5b85 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -414,8 +414,10 @@ public void restore() { // when the consumer gets fenced, all its tasks should be migrated throw new TaskMigratedException("Restore consumer get fenced by instance-id polling records.", e); } catch (final InvalidOffsetException e) { - log.warn("Encountered {} fetching records from restore consumer for partitions {}, " + - "marking the corresponding tasks as corrupted.", e.getClass().getName(), e.partitions()); + log.warn("Encountered {} fetching records from restore consumer for partitions {}, it is likely that " + + "the consumer's position has fallen out of the topic partition offset range because the topic was " + + "truncated or compacted on the broker, marking the corresponding tasks as corrupted and re-initializing" + "it later.", e.getClass().getName(), e.partitions()); final Set taskIds = new HashSet<>(); for (final TopicPartition partition : e.partitions()) { From bbd19bc104d8ad962edfde4b72eb1c1f211f2f05 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Sat, 8 Feb 2020 10:29:17 -0800 Subject: [PATCH 08/17] minor fix --- .../kafka/streams/processor/internals/StoreChangelogReader.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index 889e9332c5b85..e7924bb7ea51c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -416,7 +416,7 @@ public void restore() { } catch (final InvalidOffsetException e) { log.warn("Encountered {} fetching records from restore consumer for partitions {}, it is likely that " + "the consumer's position has fallen out of the topic partition offset range because the topic was " + - "truncated or compacted on the broker, marking the corresponding tasks as corrupted and re-initializing" + "truncated or compacted on the broker, marking the corresponding tasks as corrupted and re-initializing" + "it later.", e.getClass().getName(), e.partitions()); final Set taskIds = new HashSet<>(); From 6b6c391334df1d1461f0587ac269d605ccb2c908 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 10 Feb 2020 14:29:11 -0800 Subject: [PATCH 09/17] github comments --- .../errors/TaskCorruptedException.java | 14 ++++++----- .../processor/internals/AbstractTask.java | 13 ++++++++++ .../internals/ProcessorStateManager.java | 24 +++++++++++++++++-- .../processor/internals/StandbyTask.java | 22 ++++++++--------- .../internals/StoreChangelogReader.java | 7 +++--- .../processor/internals/StreamTask.java | 14 +++++++---- .../processor/internals/StreamThread.java | 4 ++-- .../streams/processor/internals/Task.java | 2 ++ .../processor/internals/TaskManager.java | 9 +++++-- 9 files changed, 77 insertions(+), 32 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java b/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java index e63c9c0e86e85..f1a3b48320b75 100644 --- a/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java +++ b/streams/src/main/java/org/apache/kafka/streams/errors/TaskCorruptedException.java @@ -16,8 +16,10 @@ */ package org.apache.kafka.streams.errors; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.streams.processor.TaskId; +import java.util.Map; import java.util.Set; /** @@ -30,15 +32,15 @@ */ public class TaskCorruptedException extends StreamsException { - private final Set taskIds; + private final Map> taskWithChangelogs; - public TaskCorruptedException(final Set taskIds) { - super("Tasks " + taskIds + " are corrupted and hence needs to be re-initialized"); + public TaskCorruptedException(final Map> taskWithChangelogs) { + super("Tasks with changelogs " + taskWithChangelogs + " are corrupted and hence needs to be re-initialized"); - this.taskIds = taskIds; + this.taskWithChangelogs = taskWithChangelogs; } - public Set corruptedTaskIds() { - return taskIds; + public Map> corruptedTaskWithChangelogs() { + return taskWithChangelogs; } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java index fc98117eb47a8..9c75dbbd63858 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java @@ -20,6 +20,8 @@ import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; +import java.util.Collection; +import java.util.Collections; import java.util.Set; import static org.apache.kafka.streams.processor.internals.Task.State.CLOSED; @@ -56,6 +58,17 @@ public Set inputPartitions() { return partitions; } + @Override + public Collection changelogPartitions() { + return stateMgr.changelogPartitions(); + } + + @Override + public void markChangelogAsCorrupted(final Set partitions) { + stateMgr.markChangelogAsCorrupted(partitions); + stateMgr.checkpoint(Collections.emptyMap()); + } + @Override public StateStore getStore(final String name) { return stateMgr.getStore(name); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java index ab29938030edf..2e15b673e21dc 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorStateManager.java @@ -37,6 +37,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.stream.Collectors; import static java.lang.String.format; @@ -84,11 +85,15 @@ public static class StateStoreMetadata { // update blindly with the given offset private Long offset; + // corrupted state store should not be included in checkpointing + private boolean corrupted; + private StateStoreMetadata(final StateStore stateStore) { this.stateStore = stateStore; this.restoreCallback = null; this.recordConverter = null; this.changelogPartition = null; + this.corrupted = false; this.offset = null; } @@ -282,6 +287,20 @@ Collection changelogPartitions() { return changelogOffsets().keySet(); } + void markChangelogAsCorrupted(final Set partitions) { + for (final StateStoreMetadata storeMetadata : stores.values()) { + if (partitions.contains(storeMetadata.changelogPartition)) { + storeMetadata.corrupted = true; + partitions.remove(storeMetadata.changelogPartition); + } + } + + if (!partitions.isEmpty()) { + throw new IllegalStateException("Some partitions " + partitions + " are not contained in the store list of task " + + taskId + " marking as corrupted, this is not expected"); + } + } + @Override public Map changelogOffsets() { // return the current offsets for those logged stores @@ -441,10 +460,11 @@ public void checkpoint(final Map writtenOffsets) { final Map checkpointingOffsets = new HashMap<>(); for (final StateStoreMetadata storeMetadata : stores.values()) { - // store is logged, persistent, and has a valid current offset + // store is logged, persistent, not corrupted, and has a valid current offset if (storeMetadata.changelogPartition != null && storeMetadata.stateStore.persistent() && - storeMetadata.offset != null) { + storeMetadata.offset != null && + !storeMetadata.corrupted) { checkpointingOffsets.put(storeMetadata.changelogPartition, storeMetadata.offset); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java index 70cbed66ae086..46930b405cc4e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java @@ -187,6 +187,16 @@ private void close(final boolean clean) { transitionTo(State.CLOSED); } + @Override + public boolean commitNeeded() { + return false; + } + + @Override + public Map changelogOffsets() { + return Collections.unmodifiableMap(stateMgr.changelogOffsets()); + } + @Override public void addRecords(final TopicPartition partition, final Iterable> records) { throw new IllegalStateException("Attempted to add records to task " + id() + " for invalid input partition " + partition); @@ -223,16 +233,4 @@ public String toString(final String indent) { return sb.toString(); } - - public boolean commitNeeded() { - return false; - } - - public Collection changelogPartitions() { - return stateMgr.changelogPartitions(); - } - - public Map changelogOffsets() { - return Collections.unmodifiableMap(stateMgr.changelogOffsets()); - } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index e7924bb7ea51c..f3d2ab9938f7b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -419,11 +419,12 @@ public void restore() { "truncated or compacted on the broker, marking the corresponding tasks as corrupted and re-initializing" + "it later.", e.getClass().getName(), e.partitions()); - final Set taskIds = new HashSet<>(); + final Map> taskWithCorruptedChangelogs = new HashMap<>(); for (final TopicPartition partition : e.partitions()) { - taskIds.add(changelogs.get(partition).stateManager.taskId()); + final TaskId taskId = changelogs.get(partition).stateManager.taskId(); + taskWithCorruptedChangelogs.computeIfAbsent(taskId, k -> new HashSet<>()).add(partition); } - throw new TaskCorruptedException(taskIds); + throw new TaskCorruptedException(taskWithCorruptedChangelogs); } catch (final KafkaException e) { throw new StreamsException("Restore consumer get unexpected error polling records.", e); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 54727a3bcf50e..ac6a4624c0997 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -436,6 +436,15 @@ private void close(final boolean clean) { transitionTo(State.CLOSED); } + @Override + public void markChangelogAsCorrupted(final Set partitions) { + stateMgr.markChangelogAsCorrupted(partitions); + + // only write a new checkpoint (excluding the corrupted partitions) if eos is disabled + if (eosDisabled) + stateMgr.checkpoint(Collections.emptyMap()); + } + /** * An active task is processable if its buffer contains data for all of its input * source topic partitions, or if it is enforced to be processable @@ -887,11 +896,6 @@ public boolean commitNeeded() { return commitNeeded; } - @Override - public Collection changelogPartitions() { - return stateMgr.changelogPartitions(); - } - @Override public Map changelogOffsets() { if (state() == State.RUNNING) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 74b4d6861b096..36921e4ed5b19 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -759,9 +759,9 @@ private void runLoop() { } } catch (final TaskCorruptedException e) { log.warn("Detected the states of tasks {} are corrupted. " + - "Will close the task as dirty and re-create and bootstrap from scratch.", e.corruptedTaskIds()); + "Will close the task as dirty and re-create and bootstrap from scratch.", e.corruptedTaskWithChangelogs()); - taskManager.handleCorruption(e.corruptedTaskIds()); + taskManager.handleCorruption(e.corruptedTaskWithChangelogs()); } catch (final TaskMigratedException e) { log.warn("Detected that the thread is being fenced. " + "This implies that this thread missed a rebalance and dropped out of the consumer group. " + diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java index 57fa18eb8dfb2..18d12bffb5b7b 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java @@ -174,6 +174,8 @@ enum TaskType { */ Map changelogOffsets(); + void markChangelogAsCorrupted(final Set partitions); + default Map purgableOffsets() { return Collections.emptyMap(); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index 2e5ddf79641af..6c0fac7737463 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -120,13 +120,18 @@ void handleRebalanceComplete() { rebalanceInProgress = false; } - void handleCorruption(final Set taskIds) { - for (final TaskId taskId : taskIds) { + void handleCorruption(final Map> taskWithChangelogs) { + for (final Map.Entry> entry : taskWithChangelogs.entrySet()) { + final TaskId taskId = entry.getKey(); final Task task = tasks.get(taskId); // this call is idempotent so even if the task is only CREATED we can still call it changelogReader.remove(task.changelogPartitions()); + // mark corrupted partitions to not be checkpointed, and then close the task as dirty + final Set corruptedPartitions = entry.getValue(); + task.markChangelogAsCorrupted(corruptedPartitions); + task.closeDirty(); task.revive(); From 2fd2eae63c36c86798e5d6e331869602edb63aac Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 10 Feb 2020 15:38:35 -0800 Subject: [PATCH 10/17] fix checkstyle --- .../apache/kafka/streams/processor/internals/StandbyTask.java | 1 - .../org/apache/kafka/streams/processor/internals/StreamTask.java | 1 - 2 files changed, 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java index 46930b405cc4e..6349d640b23c9 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTask.java @@ -29,7 +29,6 @@ import org.apache.kafka.streams.processor.internals.metrics.ThreadMetrics; import org.slf4j.Logger; -import java.util.Collection; import java.util.Collections; import java.util.Map; import java.util.Set; diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index ac6a4624c0997..6bb1720d29755 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -47,7 +47,6 @@ import java.io.StringWriter; import java.nio.ByteBuffer; import java.util.Base64; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; From 40731e90c7123dc0ff88b2438ebd9c8d2d66ae5b Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Wed, 12 Feb 2020 15:56:51 -0800 Subject: [PATCH 11/17] comments --- .../streams/errors/TaskMigratedException.java | 11 ++--- .../processor/internals/AbstractTask.java | 2 - .../processor/internals/StreamTask.java | 44 +++++++++---------- .../processor/internals/StreamThread.java | 13 +----- 4 files changed, 27 insertions(+), 43 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java b/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java index d60321e17b064..4a758b404737e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java +++ b/streams/src/main/java/org/apache/kafka/streams/errors/TaskMigratedException.java @@ -18,18 +18,15 @@ /** - * Indicates that one or more tasks got migrated to another thread. - * - * 1) if the task field is specified, then that single task should be cleaned up and closed as "zombie" while the - * thread can continue as normal; - * 2) if no tasks are specified (i.e. taskId == null), it means that the hosted thread has been fenced and all - * tasks are migrated, in which case the thread should rejoin the group + * Indicates that all tasks belongs to the thread have migrated to another thread. This exception can be thrown when + * the thread gets fenced (either by the consumer coordinator or by the transaction coordinator), which means it is + * no longer part of the group but a "zombie" already */ public class TaskMigratedException extends StreamsException { private final static long serialVersionUID = 1L; public TaskMigratedException(final String message, final Throwable throwable) { - super(message + "; It means all tasks belonging to this thread have been migrated", throwable); + super(message + "; It means all tasks belonging to this thread should be migrated", throwable); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java index 9c75dbbd63858..ee04b25ab84e2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AbstractTask.java @@ -21,7 +21,6 @@ import org.apache.kafka.streams.processor.TaskId; import java.util.Collection; -import java.util.Collections; import java.util.Set; import static org.apache.kafka.streams.processor.internals.Task.State.CLOSED; @@ -66,7 +65,6 @@ public Collection changelogPartitions() { @Override public void markChangelogAsCorrupted(final Set partitions) { stateMgr.markChangelogAsCorrupted(partitions); - stateMgr.checkpoint(Collections.emptyMap()); } @Override diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 240c45fa607ab..b7c4bc6239b34 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -410,13 +410,17 @@ private void close(final boolean clean) { // whenever we have successfully committed state, it is safe to checkpoint // the state as well no matter if EOS is enabled or not stateMgr.checkpoint(checkpointableOffsets()); - } else if (eosDisabled) { - // if from unclean close, then only need to flush state to make sure that when later - // closing the states, there's no records triggering any processing anymore; also swallow all caught exceptions - // However, for a _clean_ shutdown, we try to commit and checkpoint. If there are any exceptions, they become - // fatal for the "closeClean()" call, and the caller can try again with closeDirty() to complete the shutdown. + } else { try { + // if from unclean close, then only need to flush state to make sure that when later + // closing the states, there's no records triggering any processing anymore; also swallow all caught exceptions. + // However, for a _clean_ shutdown, we try to commit and checkpoint. If there are any exceptions, they become + // fatal for the "closeClean()" call, and the caller can try again with closeDirty() to complete the shutdown. stateMgr.flush(); + + // just re-write the checkpoint file without updating the store offsets, + // if there are any corrupted partitions then they will be excluded from the overwritten file + stateMgr.checkpoint(Collections.emptyMap()); } catch (final RuntimeException error) { log.debug("Ignoring flush error in unclean close.", error); } @@ -433,8 +437,9 @@ private void close(final boolean clean) { // if the latter throws and we re-close dirty which would close the state manager again. StateManagerUtil.closeStateManager(log, logPrefix, clean, stateMgr, stateDirectory); - // if EOS is enabled, we wipe out the whole state store since they are invalid to use anymore - if (!eosDisabled) { + // if EOS is enabled, we wipe out the whole state store for unclean close + // since they are invalid to use anymore + if (!clean && !eosDisabled) { StateManagerUtil.wipeStateStores(log, stateMgr); } @@ -450,20 +455,20 @@ private void close(final boolean clean) { transitionTo(State.CLOSED); } - @Override - public void markChangelogAsCorrupted(final Set partitions) { - stateMgr.markChangelogAsCorrupted(partitions); - - // only write a new checkpoint (excluding the corrupted partitions) if eos is disabled - if (eosDisabled) - stateMgr.checkpoint(Collections.emptyMap()); - } - /** * An active task is processable if its buffer contains data for all of its input * source topic partitions, or if it is enforced to be processable */ public boolean isProcessable(final long wallClockTime) { + if (state() == State.CLOSED || state() == State.CLOSING) { + // a task is only closing / closed when 1) task manager is closing, 2) a rebalance is undergoing; + // in either case we can just log it and move on without notifying the thread since the consumer + // would soon be updated to not return any records for this task anymore. + log.info("Stream task {} is already in {} state, skip adding records to it.", id(), state()); + + return false; + } + if (partitionGroup.allPartitionsBuffered()) { idleStartTime = RecordQueue.UNKNOWN; return true; @@ -710,13 +715,6 @@ private void closeRecordCollector(final boolean clean) { */ @Override public void addRecords(final TopicPartition partition, final Iterable> records) { - if (state() == State.CLOSED || state() == State.CLOSING) { - // a task is only closing / closed when 1) task manager is closing, 2) a rebalance is undergoing; - // in either case we can just log it and move on without notifying the thread since the consumer - // would soon be updated to not return any records for this task anymore. - log.info("Stream task {} is already in {} state, skip adding records to it.", id(), state()); - } - final int newQueueSize = partitionGroup.addRawRecords(partition, records); if (log.isTraceEnabled()) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 4cb9542edf744..cfa3b4d973536 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -719,20 +719,11 @@ public void run() { try { runLoop(); cleanRun = true; - } catch (final StreamsException e) { - // do not need to rethrow - log.error("Encountered the following Streams exception during processing, " + - "this indicates an unrecoverable error and the thread is going to shutdown: ", e); - throw e; - } catch (final KafkaException e) { - log.error("Encountered the following unexpected Kafka exception during processing, " + - "this indicates an internal error and thread is going to shutdown:", e); - throw e; } catch (final Exception e) { // we have caught all Kafka related exceptions, and other runtime exceptions // should be due to user application errors - log.error("Encountered the following unexpected exception during processing " + - "and the thread is going to shutdown:", e); + log.error("Encountered the following exception during processing " + + "and the thread is going to shutdown: ", e); throw e; } finally { completeShutdown(cleanRun); From 9831638e2af3d59c85253bd809d3a6114c4ec98d Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Fri, 14 Feb 2020 12:57:17 -0800 Subject: [PATCH 12/17] move javadoc --- .../streams/processor/internals/StoreChangelogReader.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index f3d2ab9938f7b..039a8f71d79fa 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -100,10 +100,6 @@ static class ChangelogMetadata { private long totalRestored; - // only for active restoring tasks (for standby changelog it is null) - // NOTE we do not book keep the current offset since we leverage state manager as its source of truth - - // the end offset beyond which records should not be applied (yet) to restore the states // // for both active restoring tasks and standby updating tasks, it is defined as: @@ -113,6 +109,8 @@ static class ChangelogMetadata { // the log-end-offset only needs to be updated once and only need to be for active tasks since for standby // tasks it would never "complete" based on the end-offset; // the committed-offset needs to be updated periodically for those standby tasks + // + // NOTE we do not book keep the current offset since we leverage state manager as its source of truth private Long restoreEndOffset; // buffer records polled by the restore consumer; From 946a944c52361e7998bc38f6664c6ad59dbb5874 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Tue, 18 Feb 2020 16:00:42 -0800 Subject: [PATCH 13/17] further fix --- .../processor/internals/StreamTask.java | 75 ++++++++++--------- .../processor/internals/TaskManager.java | 9 ++- 2 files changed, 47 insertions(+), 37 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index b7c4bc6239b34..aad5f46a57695 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -57,6 +57,7 @@ import java.util.stream.Collectors; import static java.lang.String.format; +import static java.util.Collections.emptyMap; import static java.util.Collections.singleton; import static org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency; @@ -228,25 +229,32 @@ public void suspend() { if (state() == State.CREATED || state() == State.CLOSING || state() == State.SUSPENDED) { // do nothing log.trace("Skip suspending since state is {}", state()); + } else if (state() == State.RUNNING) { + closeTopology(true); + + commitState(); + // whenever we have successfully committed state during suspension, it is safe to checkpoint + // the state as well no matter if EOS is enabled or not + stateMgr.checkpoint(checkpointableOffsets()); + + // we should also clear any buffered records of a task when suspending it + partitionGroup.clear(); + + transitionTo(State.SUSPENDED); + log.info("Suspended running"); + } else if (state() == State.RESTORING) { + // we just checkpoint the position that we've restored up to without + // going through the commit process + stateMgr.flush(); + stateMgr.checkpoint(emptyMap()); + + // we should also clear any buffered records of a task when suspending it + partitionGroup.clear(); + + transitionTo(State.SUSPENDED); + log.info("Suspended running"); } else { - if (state() == State.RUNNING) { - closeTopology(true); - } - - if (state() == State.RUNNING || state() == State.RESTORING) { - commitState(); - // whenever we have successfully committed state during suspension, it is safe to checkpoint - // the state as well no matter if EOS is enabled or not - stateMgr.checkpoint(checkpointableOffsets()); - - // we should also clear any buffered records of a task when suspending it - partitionGroup.clear(); - - transitionTo(State.SUSPENDED); - log.info("Suspended running"); - } else { - throw new IllegalStateException("Illegal state " + state() + " while suspending active task " + id); - } + throw new IllegalStateException("Illegal state " + state() + " while suspending active task " + id); } } @@ -402,30 +410,23 @@ private void close(final boolean clean) { } else { if (state() == State.RUNNING) { closeTopology(clean); - } - if (state() == State.RUNNING || state() == State.RESTORING) { if (clean) { commitState(); // whenever we have successfully committed state, it is safe to checkpoint // the state as well no matter if EOS is enabled or not stateMgr.checkpoint(checkpointableOffsets()); } else { - try { - // if from unclean close, then only need to flush state to make sure that when later - // closing the states, there's no records triggering any processing anymore; also swallow all caught exceptions. - // However, for a _clean_ shutdown, we try to commit and checkpoint. If there are any exceptions, they become - // fatal for the "closeClean()" call, and the caller can try again with closeDirty() to complete the shutdown. - stateMgr.flush(); - - // just re-write the checkpoint file without updating the store offsets, - // if there are any corrupted partitions then they will be excluded from the overwritten file - stateMgr.checkpoint(Collections.emptyMap()); - } catch (final RuntimeException error) { - log.debug("Ignoring flush error in unclean close.", error); - } + executeAndMaybeSwallow(false, stateMgr::flush, "state manager flush"); } + transitionTo(State.CLOSING); + } else if (state() == State.RESTORING) { + executeAndMaybeSwallow(clean, () -> { + stateMgr.flush(); + stateMgr.checkpoint(Collections.emptyMap()); + }, "state manager flush and checkpoint"); + transitionTo(State.CLOSING); } else if (state() == State.SUSPENDED) { // do not need to commit / checkpoint, since when suspending we've already committed the state @@ -443,7 +444,7 @@ private void close(final boolean clean) { StateManagerUtil.wipeStateStores(log, stateMgr); } - closeRecordCollector(clean); + executeAndMaybeSwallow(clean, recordCollector::close, "record collector close"); } else { throw new IllegalStateException("Illegal state " + state() + " while closing active task " + id); } @@ -696,12 +697,14 @@ private void closeTopology(final boolean clean) { } } - private void closeRecordCollector(final boolean clean) { + private void executeAndMaybeSwallow(final boolean clean, final Runnable runnable, final String name) { try { - recordCollector.close(); + runnable.run(); } catch (final RuntimeException e) { if (clean) { throw e; + } else { + log.debug("Ignoring error in unclean {}", name); } } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index 6c0fac7737463..9b8af8d048472 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -132,7 +132,14 @@ void handleCorruption(final Map> taskWithChangelogs) final Set corruptedPartitions = entry.getValue(); task.markChangelogAsCorrupted(corruptedPartitions); - task.closeDirty(); + try { + task.closeClean(); + } catch (final RuntimeException e) { + log.error("Failed to close task {} cleanly. Attempting to re-close it as dirty.", task.id()); + // We've already recorded the exception (which is the point of clean). + // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. + task.closeDirty(); + } task.revive(); } From c7719c215917908b91d287e371c8ab64b4e93282 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 20 Feb 2020 14:36:39 -0800 Subject: [PATCH 14/17] Update streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java Co-Authored-By: John Roesler --- .../apache/kafka/streams/processor/internals/StreamTask.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index b7c4bc6239b34..fe000787af2f5 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -464,7 +464,7 @@ public boolean isProcessable(final long wallClockTime) { // a task is only closing / closed when 1) task manager is closing, 2) a rebalance is undergoing; // in either case we can just log it and move on without notifying the thread since the consumer // would soon be updated to not return any records for this task anymore. - log.info("Stream task {} is already in {} state, skip adding records to it.", id(), state()); + log.info("Stream task {} is already in {} state, skip processing it.", id(), state()); return false; } From 9210ceb07077ca44b1cdac3588a79c0e5867e851 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 20 Feb 2020 14:36:52 -0800 Subject: [PATCH 15/17] Update streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java Co-Authored-By: John Roesler --- .../apache/kafka/streams/processor/internals/StreamThread.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index cfa3b4d973536..800974945bf95 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -723,7 +723,7 @@ public void run() { // we have caught all Kafka related exceptions, and other runtime exceptions // should be due to user application errors log.error("Encountered the following exception during processing " + - "and the thread is going to shutdown: ", e); + "and the thread is going to shut down: ", e); throw e; } finally { completeShutdown(cleanRun); From 5f3628829892c4a96a51cfa95aba026e808a6510 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 20 Feb 2020 14:37:05 -0800 Subject: [PATCH 16/17] Update streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java Co-Authored-By: John Roesler --- .../apache/kafka/streams/processor/internals/StreamThread.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 800974945bf95..1bca36901a853 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -743,7 +743,7 @@ private void runLoop() { try { runOnce(); if (assignmentErrorCode.get() == AssignorError.VERSION_PROBING.code()) { - log.info("Version probing detected. Rejoin the consumer group to trigger a new rebalance."); + log.info("Version probing detected. Rejoining the consumer group to trigger a new rebalance."); assignmentErrorCode.set(AssignorError.NONE.code()); enforceRebalance(); From 03f4778cba8697bad6e5c5d7ce40df6e59214c02 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 20 Feb 2020 15:03:53 -0800 Subject: [PATCH 17/17] add unit tests --- .../processor/internals/TaskManager.java | 4 +- .../processor/internals/StreamTaskTest.java | 44 ++++++++++++++++++- 2 files changed, 44 insertions(+), 4 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java index af31d7b74fdcf..48747e931ae20 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/TaskManager.java @@ -136,9 +136,7 @@ void handleCorruption(final Map> taskWithChangelogs) try { task.closeClean(); } catch (final RuntimeException e) { - log.error("Failed to close task {} cleanly. Attempting to re-close it as dirty.", task.id()); - // We've already recorded the exception (which is the point of clean). - // Now, we should go ahead and complete the close because a half-closed task is no good to anyone. + log.error("Failed to close task {} cleanly while handling corrupted tasks. Attempting to re-close it as dirty.", task.id()); task.closeDirty(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java index 9909e9c0e993c..0a4315b6b750e 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java @@ -1321,6 +1321,46 @@ public void shouldNotCommitAndThrowOnCloseDirty() { verify(stateManager); } + @Test + public void shouldNotCommitOnSuspendRestoring() { + stateManager.flush(); + EasyMock.expectLastCall(); + stateManager.checkpoint(EasyMock.eq(Collections.emptyMap())); + EasyMock.expectLastCall(); + recordCollector.commit(EasyMock.anyObject()); + EasyMock.expectLastCall().andThrow(new AssertionError("Should not call this function")).anyTimes(); + EasyMock.replay(stateManager); + + task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); + + task.initializeIfNeeded(); + task.suspend(); + + assertEquals(Task.State.SUSPENDED, task.state()); + + verify(stateManager); + } + + @Test + public void shouldNotCommitOnCloseRestoring() { + stateManager.flush(); + EasyMock.expectLastCall(); + stateManager.checkpoint(EasyMock.eq(Collections.emptyMap())); + EasyMock.expectLastCall(); + recordCollector.commit(EasyMock.anyObject()); + EasyMock.expectLastCall().andThrow(new AssertionError("Should not call this function")).anyTimes(); + EasyMock.replay(stateManager); + + task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); + + task.initializeIfNeeded(); + task.closeClean(); + + assertEquals(Task.State.CLOSED, task.state()); + + verify(stateManager); + } + @Test public void shouldCommitOnCloseClean() { final long offset = 543L; @@ -1339,6 +1379,7 @@ public void shouldCommitOnCloseClean() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); task.initializeIfNeeded(); + task.completeRestoration(); task.closeClean(); assertEquals(Task.State.CLOSED, task.state()); @@ -1365,6 +1406,7 @@ public void shouldThrowOnCloseCleanError() { task = createOptimizedStatefulTask(createConfig(false, "100"), consumer); task.initializeIfNeeded(); + task.completeRestoration(); assertThrows(ProcessorStateException.class, task::closeClean); @@ -1418,7 +1460,7 @@ public void shouldThrowOnCloseCleanCheckpointError() { EasyMock.expectLastCall(); stateManager.flush(); EasyMock.expectLastCall(); - stateManager.checkpoint(Collections.singletonMap(changelogPartition, offset)); + stateManager.checkpoint(Collections.emptyMap()); EasyMock.expectLastCall().andThrow(new ProcessorStateException("KABOOM!")).anyTimes(); stateManager.close(); EasyMock.expectLastCall().andThrow(new AssertionError("Close should not be called!")).anyTimes();