From 0d3aceb5702c0f07ebae8524ec973abd8126695a Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 9 Mar 2020 17:25:27 -0700 Subject: [PATCH 1/7] first try --- .../org/apache/kafka/common/utils/Utils.java | 29 ++++-- .../processor/internals/StateDirectory.java | 92 ++++++++++--------- .../processor/internals/TaskManager.java | 2 +- .../internals/StateDirectoryTest.java | 16 +++- .../processor/internals/TaskManagerTest.java | 2 +- 5 files changed, 84 insertions(+), 57 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 0d16c0a811c5d..61f47b4f1b963 100755 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -734,31 +734,44 @@ public static Properties mkProperties(final Map properties) { } /** - * Recursively delete the given file/directory and any subfiles (if any exist) + * Recursively delete the given file/directory and any subfiles (if any exist); + * if there are specified subfiles to keep, then maitain those files as well as the parent file * - * @param file The root file at which to begin deleting + * @param rootFile The root file at which to begin deleting + * @param filesToKeep The subfiles to keep */ - public static void delete(final File file) throws IOException { - if (file == null) + public static void delete(final File rootFile, final String... filesToKeep) throws IOException { + if (rootFile == null) return; - Files.walkFileTree(file.toPath(), new SimpleFileVisitor() { + final List files = Arrays.asList(filesToKeep); + Files.walkFileTree(rootFile.toPath(), new SimpleFileVisitor() { @Override public FileVisitResult visitFileFailed(Path path, IOException exc) throws IOException { // If the root path did not exist, ignore the error; otherwise throw it. - if (exc instanceof NoSuchFileException && path.toFile().equals(file)) + if (exc instanceof NoSuchFileException && path.toFile().equals(rootFile)) return FileVisitResult.TERMINATE; throw exc; } @Override public FileVisitResult visitFile(Path path, BasicFileAttributes attrs) throws IOException { - Files.delete(path); + if (!files.contains(path.toFile().getName())) { + Files.delete(path); + } return FileVisitResult.CONTINUE; } @Override public FileVisitResult postVisitDirectory(Path path, IOException exc) throws IOException { - Files.delete(path); + // KAFKA-8999: if there's an exception thrown previously already, we should throw it + if (exc != null) { + throw exc; + } + + // only delete the parent directory if there's nothing to keep + if (files.isEmpty()) + Files.delete(path); + return FileVisitResult.CONTINUE; } }); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java index 206867802cdd6..a7801b19f46d8 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java @@ -50,11 +50,12 @@ public class StateDirectory { static final String LOCK_FILE_NAME = ".lock"; private static final Logger log = LoggerFactory.getLogger(StateDirectory.class); + private final Time time; + private final String appId; private final File stateDir; private final boolean createStateDirectory; private final HashMap channels = new HashMap<>(); private final HashMap locks = new HashMap<>(); - private final Time time; private FileChannel globalStateChannel; private FileLock globalStateLock; @@ -75,18 +76,17 @@ private static class LockAndOwner { * @throws ProcessorStateException if the base state directory or application state directory does not exist * and could not be created when createStateDirectory is enabled. */ - public StateDirectory(final StreamsConfig config, - final Time time, - final boolean createStateDirectory) { + public StateDirectory(final StreamsConfig config, final Time time, final boolean createStateDirectory) { this.time = time; this.createStateDirectory = createStateDirectory; + this.appId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); final String stateDirName = config.getString(StreamsConfig.STATE_DIR_CONFIG); final File baseDir = new File(stateDirName); if (this.createStateDirectory && !baseDir.exists() && !baseDir.mkdirs()) { throw new ProcessorStateException( String.format("base state directory [%s] doesn't exist and couldn't be created", stateDirName)); } - stateDir = new File(baseDir, config.getString(StreamsConfig.APPLICATION_ID_CONFIG)); + stateDir = new File(baseDir, appId); if (this.createStateDirectory && !stateDir.exists() && !stateDir.mkdir()) { throw new ProcessorStateException( String.format("state directory [%s] doesn't exist and couldn't be created", stateDir.getPath())); @@ -113,9 +113,13 @@ public File directoryForTask(final TaskId taskId) { boolean directoryForTaskIsEmpty(final TaskId taskId) { final File taskDir = directoryForTask(taskId); + return taskDirEmpty(taskDir); + } + + private boolean taskDirEmpty(final File taskDir) { final File[] storeDirs = taskDir.listFiles(pathname -> !pathname.getName().equals(LOCK_FILE_NAME) && - !pathname.getName().equals(CHECKPOINT_FILE_NAME)); + !pathname.getName().equals(CHECKPOINT_FILE_NAME)); // if the task is stateless, storeDirs would be null return storeDirs == null || storeDirs.length == 0; @@ -252,18 +256,29 @@ synchronized void unlock(final TaskId taskId) throws IOException { } public synchronized void clean() { + // remove task dirs try { cleanRemovedTasks(0, true); } catch (final Exception e) { // this is already logged within cleanRemovedTasks throw new StreamsException(e); } + // remove global dir try { if (stateDir.exists()) { Utils.delete(globalStateDir().getAbsoluteFile()); } } catch (final IOException e) { - log.error("{} Failed to delete global state directory due to an unexpected exception", logPrefix(), e); + log.error("{} Failed to delete global state directory of {} due to an unexpected exception", + appId, logPrefix(), e); + throw new StreamsException(e); + } + // finally remove the parent state dir + try { + Utils.delete(stateDir); + } catch (final IOException e) { + log.error("{} Failed to delete the state directory of {} due to an unexpected exception", + appId, logPrefix(), e); throw new StreamsException(e); } } @@ -285,7 +300,7 @@ public synchronized void cleanRemovedTasks(final long cleanupDelayMs) { private synchronized void cleanRemovedTasks(final long cleanupDelayMs, final boolean manualUserCall) throws Exception { - final File[] taskDirs = listTaskDirectories(); + final File[] taskDirs = listNonEmptyTaskDirectories(); if (taskDirs == null || taskDirs.length == 0) { return; // nothing to do } @@ -294,61 +309,54 @@ private synchronized void cleanRemovedTasks(final long cleanupDelayMs, final String dirName = taskDir.getName(); final TaskId id = TaskId.parse(dirName); if (!locks.containsKey(id)) { + Exception exception = null; try { if (lock(id)) { final long now = time.milliseconds(); final long lastModifiedMs = taskDir.lastModified(); - if (now > lastModifiedMs + cleanupDelayMs || manualUserCall) { - if (!manualUserCall) { - log.info( - "{} Deleting obsolete state directory {} for task {} as {}ms has elapsed (cleanup delay is {}ms).", - logPrefix(), - dirName, - id, - now - lastModifiedMs, - cleanupDelayMs); - } else { - log.info( - "{} Deleting state directory {} for task {} as user calling cleanup.", - logPrefix(), - dirName, - id); - } + if (now > lastModifiedMs + cleanupDelayMs) { + log.info("{} Deleting obsolete state directory {} for task {} as {}ms has elapsed (cleanup delay is {}ms).", + logPrefix(), dirName, id, now - lastModifiedMs, cleanupDelayMs); + + Utils.delete(taskDir); + } else if (manualUserCall) { + log.info("{} Deleting state directory {} for task {} as user calling cleanup.", + logPrefix(), dirName, id); + Utils.delete(taskDir); } } - } catch (final OverlappingFileLockException e) { - // locked by another thread - if (manualUserCall) { - log.error("{} Failed to get the state directory lock.", logPrefix(), e); - throw e; - } - } catch (final IOException e) { - log.error("{} Failed to delete the state directory.", logPrefix(), e); - if (manualUserCall) { - throw e; - } + } catch (final OverlappingFileLockException | IOException e) { + exception = e; } finally { try { unlock(id); } catch (final IOException e) { - log.error("{} Failed to release the state directory lock.", logPrefix()); - if (manualUserCall) { - throw e; - } + exception = e; } } + + if (exception != null && manualUserCall) { + log.error("{} Failed to release the state directory lock.", logPrefix()); + throw exception; + } } } } /** - * List all of the task directories + * List all of the task directories that are non-empty * @return The list of all the existing local directories for stream tasks */ - File[] listTaskDirectories() { + File[] listNonEmptyTaskDirectories() { return !stateDir.exists() ? new File[0] : - stateDir.listFiles(pathname -> pathname.isDirectory() && PATH_NAME.matcher(pathname.getName()).matches()); + stateDir.listFiles(pathname -> { + if (!pathname.isDirectory() || !PATH_NAME.matcher(pathname.getName()).matches()) { + return false; + } else { + return !taskDirEmpty(pathname); + } + }); } private FileChannel getOrCreateFileChannel(final TaskId taskId, 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 ba75f86b0c321..c6d6784b1b2c1 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 @@ -395,7 +395,7 @@ private Set tasksOnLocalStorage() { final Set locallyStoredTasks = new HashSet<>(); - final File[] stateDirs = stateDirectory.listTaskDirectories(); + final File[] stateDirs = stateDirectory.listNonEmptyTaskDirectories(); if (stateDirs != null) { for (final File dir : stateDirs) { try { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java index e1ca918848bf4..5c351afb1f380 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java @@ -290,15 +290,21 @@ public void shouldNotRemoveNonTaskDirectoriesAndFiles() { } @Test - public void shouldListAllTaskDirectories() { + public void shouldNotListNonEmptyTaskDirectories() { TestUtils.tempDirectory(stateDir.toPath(), "foo"); final File taskDir1 = directory.directoryForTask(new TaskId(0, 0)); final File taskDir2 = directory.directoryForTask(new TaskId(0, 1)); - final List dirs = Arrays.asList(directory.listTaskDirectories()); - assertEquals(2, dirs.size()); - assertTrue(dirs.contains(taskDir1)); - assertTrue(dirs.contains(taskDir2)); + final File storeDir = new File(taskDir1, "store"); + assertTrue(storeDir.mkdir()); + + List dirs = Arrays.asList(directory.listNonEmptyTaskDirectories()); + assertEquals(Collections.singletonList(taskDir1), dirs); + + directory.cleanRemovedTasks(0L); + + dirs = Arrays.asList(directory.listNonEmptyTaskDirectories()); + assertEquals(Collections.emptyList(), dirs); } @Test 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 ef3cd55982e28..acf507b01f83b 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 @@ -168,7 +168,7 @@ public void shouldReturnOffsetsForAllCachedTaskIdsFromDirectory() throws IOExcep assertThat((new File(taskFolders[1], StateManagerUtil.CHECKPOINT_FILE_NAME)).createNewFile(), is(true)); assertThat((new File(taskFolders[3], StateManagerUtil.CHECKPOINT_FILE_NAME)).createNewFile(), is(true)); - expect(stateDirectory.listTaskDirectories()).andReturn(taskFolders).once(); + expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(taskFolders).once(); replay(activeTaskCreator, stateDirectory); From 7e0a5038991cc473f376c8d43c8a26c8a484b404 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 9 Mar 2020 17:26:23 -0700 Subject: [PATCH 2/7] pass in LOCK_FILE_NAME --- .../kafka/streams/processor/internals/StateDirectory.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java index a7801b19f46d8..4fa87d4cb81b6 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java @@ -318,12 +318,12 @@ private synchronized void cleanRemovedTasks(final long cleanupDelayMs, log.info("{} Deleting obsolete state directory {} for task {} as {}ms has elapsed (cleanup delay is {}ms).", logPrefix(), dirName, id, now - lastModifiedMs, cleanupDelayMs); - Utils.delete(taskDir); + Utils.delete(taskDir, LOCK_FILE_NAME); } else if (manualUserCall) { log.info("{} Deleting state directory {} for task {} as user calling cleanup.", logPrefix(), dirName, id); - Utils.delete(taskDir); + Utils.delete(taskDir, LOCK_FILE_NAME); } } } catch (final OverlappingFileLockException | IOException e) { From fe259171a744f9fcc52c09e008621e03d0ee0558 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 9 Mar 2020 18:15:39 -0700 Subject: [PATCH 3/7] refactoring --- .../org/apache/kafka/common/utils/Utils.java | 23 +++++++++++++------ .../processor/internals/StateDirectory.java | 5 ++-- .../internals/StateDirectoryTest.java | 3 +++ 3 files changed, 22 insertions(+), 9 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 61f47b4f1b963..3ddbdb7a13a01 100755 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -734,16 +734,25 @@ public static Properties mkProperties(final Map properties) { } /** - * Recursively delete the given file/directory and any subfiles (if any exist); - * if there are specified subfiles to keep, then maitain those files as well as the parent file + * Recursively delete the given file/directory and any subfiles (if any exist) * * @param rootFile The root file at which to begin deleting - * @param filesToKeep The subfiles to keep */ - public static void delete(final File rootFile, final String... filesToKeep) throws IOException { + public static void delete(final File rootFile) throws IOException { + delete(rootFile, Collections.emptyList()); + } + + /** + * Recursively delete the subfiles (if any exist) of the passed in root file that are not included + * in the list to keep + * + * @param rootFile The root file at which to begin deleting + * @param filesToKeep The subfiles to keep (note that if a subfile is to be kept, so are all its parent + * files in its pat)h; if empty we would also delete the root file + */ + public static void delete(final File rootFile, final List filesToKeep) throws IOException { if (rootFile == null) return; - final List files = Arrays.asList(filesToKeep); Files.walkFileTree(rootFile.toPath(), new SimpleFileVisitor() { @Override public FileVisitResult visitFileFailed(Path path, IOException exc) throws IOException { @@ -755,7 +764,7 @@ public FileVisitResult visitFileFailed(Path path, IOException exc) throws IOExce @Override public FileVisitResult visitFile(Path path, BasicFileAttributes attrs) throws IOException { - if (!files.contains(path.toFile().getName())) { + if (!filesToKeep.contains(path.toFile())) { Files.delete(path); } return FileVisitResult.CONTINUE; @@ -769,7 +778,7 @@ public FileVisitResult postVisitDirectory(Path path, IOException exc) throws IOE } // only delete the parent directory if there's nothing to keep - if (files.isEmpty()) + if (filesToKeep.isEmpty()) Files.delete(path); return FileVisitResult.CONTINUE; diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java index 4fa87d4cb81b6..65182d1714075 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java @@ -33,6 +33,7 @@ import java.nio.file.NoSuchFileException; import java.nio.file.Path; import java.nio.file.StandardOpenOption; +import java.util.Collections; import java.util.HashMap; import java.util.regex.Pattern; @@ -318,12 +319,12 @@ private synchronized void cleanRemovedTasks(final long cleanupDelayMs, log.info("{} Deleting obsolete state directory {} for task {} as {}ms has elapsed (cleanup delay is {}ms).", logPrefix(), dirName, id, now - lastModifiedMs, cleanupDelayMs); - Utils.delete(taskDir, LOCK_FILE_NAME); + Utils.delete(taskDir, Collections.singletonList(new File(taskDir, LOCK_FILE_NAME))); } else if (manualUserCall) { log.info("{} Deleting state directory {} for task {} as user calling cleanup.", logPrefix(), dirName, id); - Utils.delete(taskDir, LOCK_FILE_NAME); + Utils.delete(taskDir, Collections.singletonList(new File(taskDir, LOCK_FILE_NAME))); } } } catch (final OverlappingFileLockException | IOException e) { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java index 5c351afb1f380..1851395743d48 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java @@ -261,6 +261,9 @@ public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws directory.cleanRemovedTasks(0); files = Arrays.asList(Objects.requireNonNull(appDir.listFiles())); + assertEquals(3, files.size()); + + files = Arrays.asList(Objects.requireNonNull(directory.listNonEmptyTaskDirectories())); assertEquals(2, files.size()); assertTrue(files.contains(new File(appDir, task0.toString()))); assertTrue(files.contains(new File(appDir, task1.toString()))); From 22844d00ccdf50a003f4407d5ba167969c3014dd Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Mon, 9 Mar 2020 20:15:43 -0700 Subject: [PATCH 4/7] final fix --- .../org/apache/kafka/common/utils/Utils.java | 9 +++++-- .../processor/internals/StateDirectory.java | 27 ++++++++++++------- .../internals/StateDirectoryTest.java | 17 +++++++++--- 3 files changed, 38 insertions(+), 15 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 3ddbdb7a13a01..e9d4cc40eee34 100755 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -777,9 +777,14 @@ public FileVisitResult postVisitDirectory(Path path, IOException exc) throws IOE throw exc; } - // only delete the parent directory if there's nothing to keep - if (filesToKeep.isEmpty()) + if (rootFile.toPath().equals(path)) { + // only delete the parent directory if there's nothing to keep + if (filesToKeep.isEmpty()) { + Files.delete(path); + } + } else if (!filesToKeep.contains(path.toFile())) { Files.delete(path); + } return FileVisitResult.CONTINUE; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java index 65182d1714075..bc97c2c034f8a 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java @@ -274,14 +274,6 @@ public synchronized void clean() { appId, logPrefix(), e); throw new StreamsException(e); } - // finally remove the parent state dir - try { - Utils.delete(stateDir); - } catch (final IOException e) { - log.error("{} Failed to delete the state directory of {} due to an unexpected exception", - appId, logPrefix(), e); - throw new StreamsException(e); - } } /** @@ -301,7 +293,7 @@ public synchronized void cleanRemovedTasks(final long cleanupDelayMs) { private synchronized void cleanRemovedTasks(final long cleanupDelayMs, final boolean manualUserCall) throws Exception { - final File[] taskDirs = listNonEmptyTaskDirectories(); + final File[] taskDirs = lisAllTaskDirectories(); if (taskDirs == null || taskDirs.length == 0) { return; // nothing to do } @@ -332,6 +324,12 @@ private synchronized void cleanRemovedTasks(final long cleanupDelayMs, } finally { try { unlock(id); + + // for manual user call, stream threads are not running so it is safe to delete + // the whole directory + if (manualUserCall) { + Utils.delete(taskDir); + } } catch (final IOException e) { exception = e; } @@ -347,7 +345,7 @@ private synchronized void cleanRemovedTasks(final long cleanupDelayMs, /** * List all of the task directories that are non-empty - * @return The list of all the existing local directories for stream tasks + * @return The list of all the non-empty local directories for stream tasks */ File[] listNonEmptyTaskDirectories() { return !stateDir.exists() ? new File[0] : @@ -360,6 +358,15 @@ File[] listNonEmptyTaskDirectories() { }); } + /** + * List all of the task directories + * @return The list of all the existing local directories for stream tasks + */ + File[] lisAllTaskDirectories() { + return !stateDir.exists() ? new File[0] : + stateDir.listFiles(pathname -> pathname.isDirectory() && PATH_NAME.matcher(pathname.getName()).matches()); + } + private FileChannel getOrCreateFileChannel(final TaskId taskId, final Path lockPath) throws IOException { if (!channels.containsKey(taskId)) { diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java index 1851395743d48..a45f46e423ac5 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java @@ -249,15 +249,22 @@ public void shouldReleaseTaskStateDirectoryLock() throws Exception { public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws Exception { final TaskId task0 = new TaskId(0, 0); final TaskId task1 = new TaskId(1, 0); + final TaskId task2 = new TaskId(2, 0); try { + assertTrue(new File(directory.directoryForTask(task0), "store").mkdir()); + assertTrue(new File(directory.directoryForTask(task1), "store").mkdir()); + assertTrue(new File(directory.directoryForTask(task2), "store").mkdir()); + directory.lock(task0); directory.lock(task1); - directory.directoryForTask(new TaskId(2, 0)); List files = Arrays.asList(Objects.requireNonNull(appDir.listFiles())); assertEquals(3, files.size()); - time.sleep(1000); + files = Arrays.asList(Objects.requireNonNull(directory.listNonEmptyTaskDirectories())); + assertEquals(3, files.size()); + + time.sleep(5000); directory.cleanRemovedTasks(0); files = Arrays.asList(Objects.requireNonNull(appDir.listFiles())); @@ -276,13 +283,17 @@ public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws @Test public void shouldCleanupStateDirectoriesWhenLastModifiedIsLessThanNowMinusCleanupDelay() { final File dir = directory.directoryForTask(new TaskId(2, 0)); + assertTrue(new File(dir, "store").mkdir()); + final int cleanupDelayMs = 60000; directory.cleanRemovedTasks(cleanupDelayMs); assertTrue(dir.exists()); + assertEquals(1, directory.listNonEmptyTaskDirectories().length); time.sleep(cleanupDelayMs + 1000); directory.cleanRemovedTasks(cleanupDelayMs); - assertFalse(dir.exists()); + assertTrue(dir.exists()); + assertEquals(0, directory.listNonEmptyTaskDirectories().length); } @Test From e9618b165e97b3f54d5cff0a43ab7d812212896d Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Thu, 12 Mar 2020 18:26:55 -0700 Subject: [PATCH 5/7] unit tests --- .../processor/internals/StateDirectoryTest.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java index a45f46e423ac5..74a216a109c16 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java @@ -258,16 +258,17 @@ public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws directory.lock(task0); directory.lock(task1); - List files = Arrays.asList(Objects.requireNonNull(appDir.listFiles())); + List files = Arrays.asList(Objects.requireNonNull(directory.lisAllTaskDirectories())); assertEquals(3, files.size()); + files = Arrays.asList(Objects.requireNonNull(directory.listNonEmptyTaskDirectories())); assertEquals(3, files.size()); time.sleep(5000); directory.cleanRemovedTasks(0); - files = Arrays.asList(Objects.requireNonNull(appDir.listFiles())); + files = Arrays.asList(Objects.requireNonNull(directory.lisAllTaskDirectories())); assertEquals(3, files.size()); files = Arrays.asList(Objects.requireNonNull(directory.listNonEmptyTaskDirectories())); @@ -288,11 +289,13 @@ public void shouldCleanupStateDirectoriesWhenLastModifiedIsLessThanNowMinusClean final int cleanupDelayMs = 60000; directory.cleanRemovedTasks(cleanupDelayMs); assertTrue(dir.exists()); + assertEquals(1, directory.lisAllTaskDirectories().length); assertEquals(1, directory.listNonEmptyTaskDirectories().length); time.sleep(cleanupDelayMs + 1000); directory.cleanRemovedTasks(cleanupDelayMs); assertTrue(dir.exists()); + assertEquals(1, directory.lisAllTaskDirectories().length); assertEquals(0, directory.listNonEmptyTaskDirectories().length); } @@ -312,13 +315,13 @@ public void shouldNotListNonEmptyTaskDirectories() { final File storeDir = new File(taskDir1, "store"); assertTrue(storeDir.mkdir()); - List dirs = Arrays.asList(directory.listNonEmptyTaskDirectories()); - assertEquals(Collections.singletonList(taskDir1), dirs); + assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.lisAllTaskDirectories())); + assertEquals(Collections.singletonList(taskDir1), Arrays.asList(directory.listNonEmptyTaskDirectories())); directory.cleanRemovedTasks(0L); - dirs = Arrays.asList(directory.listNonEmptyTaskDirectories()); - assertEquals(Collections.emptyList(), dirs); + assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.lisAllTaskDirectories())); + assertEquals(Collections.emptyList(), Arrays.asList(directory.listNonEmptyTaskDirectories())); } @Test From 9e3146062c7acfbfd8aab05bc574016ad6a4a398 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Fri, 13 Mar 2020 09:45:27 -0700 Subject: [PATCH 6/7] address github comments --- .../apache/kafka/streams/KafkaStreams.java | 4 +-- .../processor/internals/StateDirectory.java | 34 +++++++++++-------- .../internals/StateDirectoryTest.java | 14 ++++---- 3 files changed, 29 insertions(+), 23 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index 45d6acb729175..2a901ed21d339 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -716,11 +716,11 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, } final ProcessorTopology globalTaskTopology = internalTopologyBuilder.buildGlobalStateTopology(); final long cacheSizePerThread = totalCacheSize / (threads.length + (globalTaskTopology == null ? 0 : 1)); - final boolean createStateDirectory = taskTopology.hasPersistentLocalStore() || + final boolean hasPersistentStores = taskTopology.hasPersistentLocalStore() || (globalTaskTopology != null && globalTaskTopology.hasPersistentGlobalStore()); try { - stateDirectory = new StateDirectory(config, time, createStateDirectory); + stateDirectory = new StateDirectory(config, time, hasPersistentStores); } catch (final ProcessorStateException fatal) { throw new StreamsException(fatal); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java index bc97c2c034f8a..896a4fd8eab84 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StateDirectory.java @@ -54,7 +54,7 @@ public class StateDirectory { private final Time time; private final String appId; private final File stateDir; - private final boolean createStateDirectory; + private final boolean hasPersistentStores; private final HashMap channels = new HashMap<>(); private final HashMap locks = new HashMap<>(); @@ -74,21 +74,27 @@ private static class LockAndOwner { /** * Ensures that the state base directory as well as the application's sub-directory are created. * + * @param config streams application configuration to read the root state directory path + * @param time system timer used to execute periodic cleanup procedure + * @param hasPersistentStores only when the application's topology does have stores persisted on local file + * system, we would go ahead and auto-create the corresponding application / task / store + * directories whenever necessary; otherwise no directories would be created. + * * @throws ProcessorStateException if the base state directory or application state directory does not exist - * and could not be created when createStateDirectory is enabled. + * and could not be created when hasPersistentStores is enabled. */ - public StateDirectory(final StreamsConfig config, final Time time, final boolean createStateDirectory) { + public StateDirectory(final StreamsConfig config, final Time time, final boolean hasPersistentStores) { this.time = time; - this.createStateDirectory = createStateDirectory; + this.hasPersistentStores = hasPersistentStores; this.appId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); final String stateDirName = config.getString(StreamsConfig.STATE_DIR_CONFIG); final File baseDir = new File(stateDirName); - if (this.createStateDirectory && !baseDir.exists() && !baseDir.mkdirs()) { + if (this.hasPersistentStores && !baseDir.exists() && !baseDir.mkdirs()) { throw new ProcessorStateException( String.format("base state directory [%s] doesn't exist and couldn't be created", stateDirName)); } stateDir = new File(baseDir, appId); - if (this.createStateDirectory && !stateDir.exists() && !stateDir.mkdir()) { + if (this.hasPersistentStores && !stateDir.exists() && !stateDir.mkdir()) { throw new ProcessorStateException( String.format("state directory [%s] doesn't exist and couldn't be created", stateDir.getPath())); } @@ -101,7 +107,7 @@ public StateDirectory(final StreamsConfig config, final Time time, final boolean */ public File directoryForTask(final TaskId taskId) { final File taskDir = new File(stateDir, taskId.toString()); - if (createStateDirectory && !taskDir.exists() && !taskDir.mkdir()) { + if (hasPersistentStores && !taskDir.exists() && !taskDir.mkdir()) { throw new ProcessorStateException( String.format("task directory [%s] doesn't exist and couldn't be created", taskDir.getPath())); } @@ -133,7 +139,7 @@ private boolean taskDirEmpty(final File taskDir) { */ File globalStateDir() { final File dir = new File(stateDir, "global"); - if (createStateDirectory && !dir.exists() && !dir.mkdir()) { + if (hasPersistentStores && !dir.exists() && !dir.mkdir()) { throw new ProcessorStateException( String.format("global state directory [%s] doesn't exist and couldn't be created", dir.getPath())); } @@ -146,12 +152,12 @@ private String logPrefix() { /** * Get the lock for the {@link TaskId}s directory if it is available - * @param taskId + * @param taskId task id * @return true if successful - * @throws IOException + * @throws IOException if the file cannot be created or file handle cannot be grabbed, should be considered as fatal */ synchronized boolean lock(final TaskId taskId) throws IOException { - if (!createStateDirectory) { + if (!hasPersistentStores) { return true; } @@ -195,7 +201,7 @@ synchronized boolean lock(final TaskId taskId) throws IOException { } synchronized boolean lockGlobalState() throws IOException { - if (!createStateDirectory) { + if (!hasPersistentStores) { return true; } @@ -293,7 +299,7 @@ public synchronized void cleanRemovedTasks(final long cleanupDelayMs) { private synchronized void cleanRemovedTasks(final long cleanupDelayMs, final boolean manualUserCall) throws Exception { - final File[] taskDirs = lisAllTaskDirectories(); + final File[] taskDirs = listAllTaskDirectories(); if (taskDirs == null || taskDirs.length == 0) { return; // nothing to do } @@ -362,7 +368,7 @@ File[] listNonEmptyTaskDirectories() { * List all of the task directories * @return The list of all the existing local directories for stream tasks */ - File[] lisAllTaskDirectories() { + File[] listAllTaskDirectories() { return !stateDir.exists() ? new File[0] : stateDir.listFiles(pathname -> pathname.isDirectory() && PATH_NAME.matcher(pathname.getName()).matches()); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java index 74a216a109c16..827557a88faae 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateDirectoryTest.java @@ -258,7 +258,7 @@ public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws directory.lock(task0); directory.lock(task1); - List files = Arrays.asList(Objects.requireNonNull(directory.lisAllTaskDirectories())); + List files = Arrays.asList(Objects.requireNonNull(directory.listAllTaskDirectories())); assertEquals(3, files.size()); @@ -268,7 +268,7 @@ public void shouldCleanUpTaskStateDirectoriesThatAreNotCurrentlyLocked() throws time.sleep(5000); directory.cleanRemovedTasks(0); - files = Arrays.asList(Objects.requireNonNull(directory.lisAllTaskDirectories())); + files = Arrays.asList(Objects.requireNonNull(directory.listAllTaskDirectories())); assertEquals(3, files.size()); files = Arrays.asList(Objects.requireNonNull(directory.listNonEmptyTaskDirectories())); @@ -289,13 +289,13 @@ public void shouldCleanupStateDirectoriesWhenLastModifiedIsLessThanNowMinusClean final int cleanupDelayMs = 60000; directory.cleanRemovedTasks(cleanupDelayMs); assertTrue(dir.exists()); - assertEquals(1, directory.lisAllTaskDirectories().length); + assertEquals(1, directory.listAllTaskDirectories().length); assertEquals(1, directory.listNonEmptyTaskDirectories().length); time.sleep(cleanupDelayMs + 1000); directory.cleanRemovedTasks(cleanupDelayMs); assertTrue(dir.exists()); - assertEquals(1, directory.lisAllTaskDirectories().length); + assertEquals(1, directory.listAllTaskDirectories().length); assertEquals(0, directory.listNonEmptyTaskDirectories().length); } @@ -307,7 +307,7 @@ public void shouldNotRemoveNonTaskDirectoriesAndFiles() { } @Test - public void shouldNotListNonEmptyTaskDirectories() { + public void shouldOnlyListNonEmptyTaskDirectories() { TestUtils.tempDirectory(stateDir.toPath(), "foo"); final File taskDir1 = directory.directoryForTask(new TaskId(0, 0)); final File taskDir2 = directory.directoryForTask(new TaskId(0, 1)); @@ -315,12 +315,12 @@ public void shouldNotListNonEmptyTaskDirectories() { final File storeDir = new File(taskDir1, "store"); assertTrue(storeDir.mkdir()); - assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.lisAllTaskDirectories())); + assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.listAllTaskDirectories())); assertEquals(Collections.singletonList(taskDir1), Arrays.asList(directory.listNonEmptyTaskDirectories())); directory.cleanRemovedTasks(0L); - assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.lisAllTaskDirectories())); + assertEquals(Arrays.asList(taskDir1, taskDir2), Arrays.asList(directory.listAllTaskDirectories())); assertEquals(Collections.emptyList(), Arrays.asList(directory.listNonEmptyTaskDirectories())); } From e76956090f1b8010efb5b1be9dc7a24a7b9884a6 Mon Sep 17 00:00:00 2001 From: Guozhang Wang Date: Sat, 14 Mar 2020 13:46:51 -0700 Subject: [PATCH 7/7] fix unit tests and renaming --- .../kafka/streams/processor/internals/TaskManager.java | 10 +++++----- .../streams/processor/internals/TaskManagerTest.java | 6 +++--- 2 files changed, 8 insertions(+), 8 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 95b641bf8791d..b804a4a7d6429 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 @@ -127,7 +127,7 @@ boolean isRebalanceInProgress() { void handleRebalanceStart(final Set subscribedTopics) { builder.addSubscribedTopicsFromMetadata(subscribedTopics, logPrefix); - tryToLockAllTaskDirectories(); + tryToLockAllNonEmptyTaskDirectories(); rebalanceInProgress = true; } @@ -379,7 +379,7 @@ void handleLostAll() { /** * Compute the offset total summed across all stores in a task. Includes offset sum for any tasks we own the - * lock for, which includes assigned and unassigned tasks we locked in {@link #tryToLockAllTaskDirectories()} + * lock for, which includes assigned and unassigned tasks we locked in {@link #tryToLockAllNonEmptyTaskDirectories()} * * @return Map from task id to its total offset summed across all state stores */ @@ -410,12 +410,12 @@ public Map getTaskOffsetSums() { } /** - * Makes a weak attempt to lock all task directories in the state dir. We are responsible for computing and + * Makes a weak attempt to lock all non-empty task directories in the state dir. We are responsible for computing and * reporting the offset sum for any unassigned tasks we obtain the lock for in the upcoming rebalance. Tasks * that we locked but didn't own will be released at the end of the rebalance (unless of course we were * assigned the task as a result of the rebalance). This method should be idempotent. */ - private void tryToLockAllTaskDirectories() { + private void tryToLockAllNonEmptyTaskDirectories() { for (final File dir : stateDirectory.listNonEmptyTaskDirectories()) { try { final TaskId id = TaskId.parse(dir.getName()); @@ -438,7 +438,7 @@ private void tryToLockAllTaskDirectories() { /** * We must release the lock for any unassigned tasks that we temporarily locked in preparation for a - * rebalance in {@link #tryToLockAllTaskDirectories()}. + * rebalance in {@link #tryToLockAllNonEmptyTaskDirectories()}. */ private void releaseLockedUnassignedTaskDirectories() { final AtomicReference firstException = new AtomicReference<>(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 1f37bb4f9ea02..2f7681d406b9d 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 @@ -168,7 +168,7 @@ public void shouldIdempotentlyUpdateSubscriptionFromActiveAssignment() { @Test public void shouldNotLockAnythingIfStateDirIsEmpty() { - expect(stateDirectory.listAllTaskDirectories()).andReturn(new File[0]).once(); + expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(new File[0]).once(); replay(stateDirectory); taskManager.handleRebalanceStart(singleton("topic")); @@ -999,7 +999,7 @@ public void shouldHandleRebalanceEvents() { expect(consumer.assignment()).andReturn(assignment); consumer.pause(assignment); expectLastCall(); - expect(stateDirectory.listAllTaskDirectories()).andReturn(new File[0]); + expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(new File[0]); replay(consumer, stateDirectory); assertThat(taskManager.isRebalanceInProgress(), is(false)); taskManager.handleRebalanceStart(emptySet()); @@ -1665,7 +1665,7 @@ private void makeTaskFolders(final String... names) throws IOException { for (int i = 0; i < names.length; ++i) { taskFolders[i] = testFolder.newFolder(names[i]); } - expect(stateDirectory.listAllTaskDirectories()).andReturn(taskFolders).once(); + expect(stateDirectory.listNonEmptyTaskDirectories()).andReturn(taskFolders).once(); } private void writeCheckpointFile(final TaskId task, final Map offsets) throws IOException {