From 682c27f08da915b620743125e6d3e9ba9fddba7c Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Mon, 13 Apr 2020 19:00:29 -0700 Subject: [PATCH 1/3] need to ensure lock --- .../processor/internals/StateManagerUtil.java | 57 ++++++++++--------- .../processor/internals/TaskManager.java | 5 -- 2 files changed, 30 insertions(+), 32 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java index 15662bf6237df..247f380be64f9 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java @@ -100,41 +100,44 @@ static void closeStateManager(final Logger log, throw new IllegalArgumentException("State store could not be wiped out during clean close"); } - ProcessorStateException exception = null; - final TaskId id = stateMgr.taskId(); - log.trace("Closing state manager for {}", id); + log.trace("Closing state manager for {} task {}", taskType, id); + ProcessorStateException exception = null; try { - stateMgr.close(); - - if (wipeStateStore) { - // we can just delete the whole dir of the task, including the state store images and the checkpoint files, - // and then we write an empty checkpoint file indicating that the previous close is graceful and we just - // need to re-bootstrap the restoration from the beginning - Utils.delete(stateMgr.baseDir()); - } - } catch (final ProcessorStateException e) { - exception = e; - } catch (final IOException e) { - throw new ProcessorStateException("Failed to wiping state stores for task " + id, e); - } finally { - try { - stateDirectory.unlock(id); - } catch (final IOException e) { - if (exception == null) { - exception = new ProcessorStateException( - String.format("%sFailed to release state dir lock", logPrefix), e); + if (stateDirectory.lock(id)) { + try { + stateMgr.close(); + + if (wipeStateStore) { + // we can just delete the whole dir of the task, including the state store images and the checkpoint files, + // and then we write an empty checkpoint file indicating that the previous close is graceful and we just + // need to re-bootstrap the restoration from the beginning + Utils.delete(stateMgr.baseDir()); + } + } catch (final ProcessorStateException e) { + exception = e; + } catch (final IOException e) { + throw new ProcessorStateException("Failed to wiping state stores for task " + id, e); + } finally { + try { + stateDirectory.unlock(id); + } catch (final IOException e) { + if (exception == null) { + exception = new ProcessorStateException(String.format("%sFailed to release state dir lock", logPrefix), e); + } + } } } + } catch (final IOException e) { + throw new StreamsException( + String.format("%sFatal error while trying to lock the state directory for task %s", logPrefix, id), + e + ); } if (exception != null) { - if (closeClean) { - throw exception; - } else { - log.warn("Closing {} task {} uncleanly and swallows an exception", taskType, id, exception); - } + throw exception; } } } 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 386be222c66a6..028ef86662a1a 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 @@ -459,7 +459,6 @@ 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()) { @@ -471,10 +470,6 @@ void handleLostAll() { log.warn("Error closing task producer for " + task.id() + " while handling lostAll", e); } } - - for (final TopicPartition inputPartition : inputPartitions) { - partitionToTask.remove(inputPartition); - } } if (processingMode == EXACTLY_ONCE_BETA) { From dad39b7e03f0098a4f4e4620b388e87770452bed Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Tue, 14 Apr 2020 16:06:56 -0700 Subject: [PATCH 2/3] fix up tests, add new ones --- .../processor/internals/StandbyTask.java | 5 +- .../processor/internals/StateManagerUtil.java | 31 +++---- .../processor/internals/StreamTask.java | 8 +- .../internals/StateManagerUtilTest.java | 85 ++++++++++++++----- .../processor/internals/StreamTaskTest.java | 3 + 5 files changed, 82 insertions(+), 50 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 aecbe2fa67caf..c44a0f99e4521 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 @@ -209,15 +209,12 @@ private void close(final boolean clean) { stateMgr.checkpoint(Collections.emptyMap()); offsetSnapshotSinceLastCommit = new HashMap<>(stateMgr.changelogOffsets()); } - final boolean wipeStateStore = !clean && eosEnabled; - log.info("standby task clean {}, eos enabled {}", clean, eosEnabled); - executeAndMaybeSwallow(clean, () -> StateManagerUtil.closeStateManager( log, logPrefix, clean, - wipeStateStore, + eosEnabled, stateMgr, stateDirectory, TaskType.STANDBY), diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java index 247f380be64f9..4813f4243628c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateManagerUtil.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.processor.internals; +import java.util.concurrent.atomic.AtomicReference; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.errors.LockException; import org.apache.kafka.streams.errors.ProcessorStateException; @@ -92,50 +93,44 @@ static void registerStateStores(final Logger log, static void closeStateManager(final Logger log, final String logPrefix, final boolean closeClean, - final boolean wipeStateStore, + final boolean eosEnabled, final ProcessorStateManager stateMgr, final StateDirectory stateDirectory, final TaskType taskType) { - if (closeClean && wipeStateStore) { - throw new IllegalArgumentException("State store could not be wiped out during clean close"); - } + // if EOS is enabled, wipe out the whole state store for unclean close since it is now invalid + final boolean wipeStateStore = !closeClean && eosEnabled; final TaskId id = stateMgr.taskId(); log.trace("Closing state manager for {} task {}", taskType, id); - ProcessorStateException exception = null; + final AtomicReference firstException = new AtomicReference<>(null); try { if (stateDirectory.lock(id)) { try { stateMgr.close(); if (wipeStateStore) { + log.debug("Wiping state stores for {} task {}", taskType, id); // we can just delete the whole dir of the task, including the state store images and the checkpoint files, // and then we write an empty checkpoint file indicating that the previous close is graceful and we just // need to re-bootstrap the restoration from the beginning Utils.delete(stateMgr.baseDir()); } } catch (final ProcessorStateException e) { - exception = e; - } catch (final IOException e) { - throw new ProcessorStateException("Failed to wiping state stores for task " + id, e); + firstException.compareAndSet(null, e); } finally { - try { - stateDirectory.unlock(id); - } catch (final IOException e) { - if (exception == null) { - exception = new ProcessorStateException(String.format("%sFailed to release state dir lock", logPrefix), e); - } - } + stateDirectory.unlock(id); } } } catch (final IOException e) { - throw new StreamsException( - String.format("%sFatal error while trying to lock the state directory for task %s", logPrefix, id), - e + final ProcessorStateException exception = new ProcessorStateException( + String.format("%sFatal error while trying to close the state manager for task %s", logPrefix, id), e ); + firstException.compareAndSet(null, exception); + } + final ProcessorStateException exception = firstException.get(); if (exception != null) { throw exception; } 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 5ad6d3d3c7e1f..51604e7d16085 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 @@ -492,11 +492,7 @@ private void close(final boolean clean, case RUNNING: case RESTORING: case SUSPENDED: - // if EOS is enabled, we wipe out the whole state store for unclean close - // since they are invalid to use anymore - final boolean wipeStateStore = !clean && eosEnabled; - - // first close state manager (which is idempotent) then close the record collector (which could throw), + // first close state manager (which is idempotent) then close the record collector // if the latter throws and we re-close dirty which would close the state manager again. executeAndMaybeSwallow( clean, @@ -504,7 +500,7 @@ private void close(final boolean clean, log, logPrefix, clean, - wipeStateStore, + eosEnabled, stateMgr, stateDirectory, TaskType.ACTIVE diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java index aa6bff961953f..f1b8d8baad490 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java @@ -23,7 +23,6 @@ import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.TaskId; import org.apache.kafka.streams.processor.internals.Task.TaskType; -import org.apache.kafka.streams.processor.internals.testutil.LogCaptureAppender; import org.apache.kafka.test.MockKeyValueStore; import org.apache.kafka.test.TestUtils; import org.easymock.IMocksControl; @@ -39,7 +38,6 @@ import java.io.File; import java.io.IOException; import java.util.Arrays; -import java.util.List; import static java.util.Collections.emptyList; import static java.util.Collections.singletonList; @@ -48,7 +46,6 @@ import static org.easymock.EasyMock.expectLastCall; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertThrows; -import static org.junit.Assert.assertTrue; import static org.powermock.api.easymock.PowerMock.mockStatic; import static org.powermock.api.easymock.PowerMock.replayAll; @@ -178,16 +175,12 @@ public void testRegisterStateStores() throws IOException { ctrl.verify(); } - @Test - public void testShouldThrowWhenCleanAndWipeStateAreBothTrue() { - assertThrows(IllegalArgumentException.class, () -> StateManagerUtil.closeStateManager(logger, - "logPrefix:", true, true, stateManager, stateDirectory, TaskType.ACTIVE)); - } - @Test public void testCloseStateManagerClean() throws IOException { expect(stateManager.taskId()).andReturn(taskId); + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); expectLastCall(); @@ -206,6 +199,9 @@ public void testCloseStateManagerClean() throws IOException { @Test public void testCloseStateManagerThrowsExceptionWhenClean() throws IOException { expect(stateManager.taskId()).andReturn(taskId); + + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); expectLastCall(); @@ -219,7 +215,6 @@ public void testCloseStateManagerThrowsExceptionWhenClean() throws IOException { ProcessorStateException.class, () -> StateManagerUtil.closeStateManager(logger, "logPrefix:", true, false, stateManager, stateDirectory, TaskType.ACTIVE)); - assertEquals("logPrefix:Failed to release state dir lock", thrown.getMessage()); assertEquals(IOException.class, thrown.getCause().getClass()); ctrl.verify(); @@ -228,6 +223,9 @@ public void testCloseStateManagerThrowsExceptionWhenClean() throws IOException { @Test public void testCloseStateManagerOnlyThrowsFirstExceptionWhenClean() throws IOException { expect(stateManager.taskId()).andReturn(taskId); + + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); expectLastCall().andThrow(new ProcessorStateException("state manager failed to close")); @@ -249,32 +247,35 @@ public void testCloseStateManagerOnlyThrowsFirstExceptionWhenClean() throws IOEx } @Test - public void testCloseStateManagerDirtyShallSwallowException() throws IOException { - final LogCaptureAppender appender = LogCaptureAppender.createAndRegister(); - + public void testCloseStateManagerThrowsExceptionWhenDirty() throws IOException { expect(stateManager.taskId()).andReturn(taskId); + + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); - expectLastCall().andThrow(new ProcessorStateException("state manager failed to close")); + expectLastCall(); stateDirectory.unlock(taskId); - expectLastCall(); + expectLastCall().andThrow(new IOException("Timeout")); ctrl.checkOrder(true); ctrl.replay(); - StateManagerUtil.closeStateManager(logger, - "logPrefix:", false, false, stateManager, stateDirectory, TaskType.ACTIVE); + final ProcessorStateException thrown = assertThrows( + ProcessorStateException.class, + () -> StateManagerUtil.closeStateManager( + logger, "logPrefix:", false, false, stateManager, stateDirectory, TaskType.ACTIVE)); - ctrl.verify(); + assertEquals(IOException.class, thrown.getCause().getClass()); - LogCaptureAppender.unregister(appender); - final List strings = appender.getMessages(); - assertTrue(strings.contains("testClosing ACTIVE task 0_0 uncleanly and swallows an exception")); + ctrl.verify(); } @Test public void testCloseStateManagerWithStateStoreWipeOut() throws IOException { expect(stateManager.taskId()).andReturn(taskId); + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); expectLastCall(); @@ -299,6 +300,8 @@ public void testCloseStateManagerWithStateStoreWipeOutRethrowWrappedIOException( mockStatic(Utils.class); expect(stateManager.taskId()).andReturn(taskId); + expect(stateDirectory.lock(taskId)).andReturn(true); + stateManager.close(); expectLastCall(); @@ -319,9 +322,47 @@ public void testCloseStateManagerWithStateStoreWipeOutRethrowWrappedIOException( ProcessorStateException.class, () -> StateManagerUtil.closeStateManager(logger, "logPrefix:", false, true, stateManager, stateDirectory, TaskType.ACTIVE)); - assertEquals("Failed to wiping state stores for task 0_0", thrown.getMessage()); assertEquals(IOException.class, thrown.getCause().getClass()); ctrl.verify(); } + + @Test + public void shouldNotStateManagerIfUnableToLockTaskDirectory() throws IOException { + expect(stateManager.taskId()).andReturn(taskId); + + expect(stateDirectory.lock(taskId)).andReturn(false); + + stateManager.close(); + expectLastCall().andThrow(new StreamsException("Should not be trying to close state you don't own!")); + + ctrl.checkOrder(true); + ctrl.replay(); + + replayAll(); + + StateManagerUtil.closeStateManager( + logger, "logPrefix:", false, true, stateManager, stateDirectory, TaskType.ACTIVE); + } + + @Test + public void shouldNotWipeStateStoresIfUnableToLockTaskDirectory() throws IOException { + final File unknownFile = new File("/unknown/path"); + expect(stateManager.taskId()).andReturn(taskId); + + expect(stateDirectory.lock(taskId)).andReturn(false); + + expect(stateManager.baseDir()).andReturn(unknownFile); + + Utils.delete(unknownFile); + expectLastCall().andThrow(new StreamsException("Should not be trying to wipe state you don't own!")); + + ctrl.checkOrder(true); + ctrl.replay(); + + replayAll(); + + StateManagerUtil.closeStateManager( + logger, "logPrefix:", false, true, stateManager, stateDirectory, TaskType.ACTIVE); + } } 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 34663707f573b..e207502335eeb 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 @@ -253,6 +253,9 @@ public void shouldAttemptToDeleteStateDirectoryWhenCloseDirtyAndEosEnabled() thr EasyMock.expect(stateManager.taskId()).andReturn(taskId); + EasyMock.expect(stateDirectory.lock(taskId)).andReturn(true); + EasyMock.expectLastCall(); + stateManager.close(); EasyMock.expectLastCall(); From 1ad8282dd07d2d4b91647f7f94fe235536139227 Mon Sep 17 00:00:00 2001 From: ableegoldman Date: Wed, 15 Apr 2020 11:28:24 -0700 Subject: [PATCH 3/3] github review --- .../streams/processor/internals/StateManagerUtilTest.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java index f1b8d8baad490..c5de83837509d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateManagerUtilTest.java @@ -328,13 +328,13 @@ public void testCloseStateManagerWithStateStoreWipeOutRethrowWrappedIOException( } @Test - public void shouldNotStateManagerIfUnableToLockTaskDirectory() throws IOException { + public void shouldNotCloseStateManagerIfUnableToLockTaskDirectory() throws IOException { expect(stateManager.taskId()).andReturn(taskId); expect(stateDirectory.lock(taskId)).andReturn(false); stateManager.close(); - expectLastCall().andThrow(new StreamsException("Should not be trying to close state you don't own!")); + expectLastCall().andThrow(new AssertionError("Should not be trying to close state you don't own!")); ctrl.checkOrder(true); ctrl.replay(); @@ -342,7 +342,7 @@ public void shouldNotStateManagerIfUnableToLockTaskDirectory() throws IOExceptio replayAll(); StateManagerUtil.closeStateManager( - logger, "logPrefix:", false, true, stateManager, stateDirectory, TaskType.ACTIVE); + logger, "logPrefix:", true, false, stateManager, stateDirectory, TaskType.ACTIVE); } @Test @@ -355,7 +355,7 @@ public void shouldNotWipeStateStoresIfUnableToLockTaskDirectory() throws IOExcep expect(stateManager.baseDir()).andReturn(unknownFile); Utils.delete(unknownFile); - expectLastCall().andThrow(new StreamsException("Should not be trying to wipe state you don't own!")); + expectLastCall().andThrow(new AssertionError("Should not be trying to wipe state you don't own!")); ctrl.checkOrder(true); ctrl.replay();