This is an automated email from the ASF dual-hosted git repository.
bbejeck pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 65d65546d55 KAFKA-20805: Refresh task directory mtime on state manager
close (#22837)
65d65546d55 is described below
commit 65d65546d554a4d1317c22b93fd9f82d843281c9
Author: Nick Telford <[email protected]>
AuthorDate: Wed Jul 15 21:00:10 2026 +0100
KAFKA-20805: Refresh task directory mtime on state manager close (#22837)
`OffsetOutOfRangeException` was observed during state restore in
long-running soak tests. A task's state directory was deleted by the
state directory cleaner roughly one second after the owning thread died
and released the directory's lock, even though the task had been
committing right up until then. When the task was reassigned it found an
empty directory, defaulted to the changelog start offset, and a
seek-to-beginning restore of a retention-trimmed changelog raced
retention and produced an out-of-range fetch.
The cleaner (`StateDirectory#cleanRemovedTasks`) deletes an unowned task
directory once its filesystem last-modified time is older than
`state.cleanup.delay.ms`. Now that changelog offsets are managed by the
state stores rather than a per-task `.checkpoint` file, nothing rewrites
a direct child of the task directory on commit, so an
actively-committing task's directory mtime no longer advances — it stays
frozen at creation time. The moment the task's lock is released, the
directory looks obsolete to the cleaner and becomes eligible for
immediate deletion.
This refreshes the task directory's last-modified time when the state
manager closes, marking the point at which the task is released. The
cleaner's grace period is then measured from release, so a just-released
directory survives long enough to be reassigned. The touch is
best-effort (`File#setLastModified` can fail on some filesystems) and
only bumps a directory that still exists.
Closing the state manager is the natural place for this: it runs on
every genuine task release (clean and dirty close, active and standby),
is never invoked by the cleaner itself, and already holds a reference to
the task directory. To keep the timestamp testable, a `Time` instance is
injected into `ProcessorStateManager` (and threaded through the task
creators) rather than reading the system clock directly.
Reviewers: Bill Bejeck <[email protected]>, Matthias J. Sax
<[email protected]>
---
.../processor/internals/ActiveTaskCreator.java | 1 +
.../processor/internals/ProcessorStateManager.java | 22 ++++++++-
.../processor/internals/StandbyTaskCreator.java | 5 ++
.../processor/internals/StateDirectory.java | 1 +
.../streams/processor/internals/StreamThread.java | 1 +
.../internals/ProcessorStateManagerTest.java | 53 +++++++++++++++++++++-
.../processor/internals/StreamThreadTest.java | 1 +
.../StreamThreadStateStoreProviderTest.java | 1 +
.../apache/kafka/streams/TopologyTestDriver.java | 1 +
9 files changed, 83 insertions(+), 3 deletions(-)
diff --git
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ActiveTaskCreator.java
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ActiveTaskCreator.java
index 5e1a69d29e0..ce3786b7bb8 100644
---
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/ActiveTaskCreator.java
+++
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/ActiveTaskCreator.java
@@ -152,6 +152,7 @@ class ActiveTaskCreator {
applicationConfig.getBoolean(StreamsConfig.TRANSACTIONAL_STATE_STORES_CONFIG),
logContext,
stateDirectory,
+ time,
topology.storeToChangelogTopic(),
partitions,
upgradeFrom);
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 6d620a56a52..501048c82d3 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
@@ -18,6 +18,7 @@ package org.apache.kafka.streams.processor.internals;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.internals.FixedOrderMap;
import org.apache.kafka.common.utils.internals.LogContext;
import org.apache.kafka.streams.errors.ProcessorStateException;
@@ -207,6 +208,9 @@ public class ProcessorStateManager implements StateManager {
private final StateDirectory stateDirectory;
private final File baseDir;
private final UpgradeFromValues upgradeFrom;
+ private final Time time;
+
+ private boolean startupTask = false;
private TaskType taskType;
private Logger log;
@@ -229,6 +233,7 @@ public class ProcessorStateManager implements StateManager {
final boolean transactionalStateStoresEnabled,
final LogContext logContext,
final StateDirectory stateDirectory,
+ final Time time,
final Map<String, String>
storeToChangelogTopic,
final Collection<TopicPartition>
sourcePartitions,
final UpgradeFromValues upgradeFrom) throws
ProcessorStateException {
@@ -241,6 +246,7 @@ public class ProcessorStateManager implements StateManager {
this.transactionalStateStoresEnabled = transactionalStateStoresEnabled;
this.sourcePartitions = sourcePartitions;
this.upgradeFrom = upgradeFrom;
+ this.time = time;
this.baseDir = stateDirectory.getOrCreateDirectoryForTask(taskId);
this.stateDirectory = stateDirectory;
@@ -259,9 +265,10 @@ public class ProcessorStateManager implements StateManager
{
final boolean transactionalStateStoresEnabled,
final LogContext logContext,
final StateDirectory stateDirectory,
+ final Time time,
final Map<String, String>
storeToChangelogTopic,
final Collection<TopicPartition>
sourcePartitions) throws ProcessorStateException {
- this(taskId, taskType, eosEnabled, transactionalStateStoresEnabled,
logContext, stateDirectory, storeToChangelogTopic, sourcePartitions, null);
+ this(taskId, taskType, eosEnabled, transactionalStateStoresEnabled,
logContext, stateDirectory, time, storeToChangelogTopic, sourcePartitions,
null);
}
/**
@@ -273,9 +280,13 @@ public class ProcessorStateManager implements StateManager
{
final boolean
eosEnabled,
final
LogContext logContext,
final
StateDirectory stateDirectory,
+ final Time time,
final
Map<String, String> storeToChangelogTopic,
final
Set<TopicPartition> sourcePartitions) {
- return new ProcessorStateManager(taskId, TaskType.STANDBY, eosEnabled,
false, logContext, stateDirectory, storeToChangelogTopic, sourcePartitions);
+ final ProcessorStateManager stateManager =
+ new ProcessorStateManager(taskId, TaskType.STANDBY, eosEnabled,
false, logContext, stateDirectory, time, storeToChangelogTopic,
sourcePartitions);
+ stateManager.startupTask = true;
+ return stateManager;
}
void registerStateStores(final List<StateStore> allStores, final
InternalProcessorContext<?, ?> processorContext) {
@@ -684,6 +695,13 @@ public class ProcessorStateManager implements StateManager
{
}
stores.clear();
+
+ // Refresh the task directory's modification time so the state
directory cleaner does
+ // not treat the just-released directory as obsolete before it can
be reassigned.
+ if (!startupTask && baseDir.exists() &&
!baseDir.setLastModified(time.milliseconds())) {
+ log.warn("{}Failed to update modification time of state
directory {} for task {}",
+ logPrefix, baseDir, taskId);
+ }
}
LegacyCheckpointingStateStore.maybeDowngradeOffsets(logPrefix,
upgradeFrom, stateDirectory, taskId, allOffsets);
diff --git
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTaskCreator.java
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTaskCreator.java
index 9083ff054e2..8afd67c5dad 100644
---
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTaskCreator.java
+++
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StandbyTaskCreator.java
@@ -18,6 +18,7 @@ package org.apache.kafka.streams.processor.internals;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.metrics.Sensor;
+import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.internals.LogContext;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.internals.UpgradeFromValues;
@@ -41,6 +42,7 @@ class StandbyTaskCreator {
private final StreamsConfig applicationConfig;
private final StreamsMetricsImpl streamsMetrics;
private final StateDirectory stateDirectory;
+ private final Time time;
private final ThreadCache dummyCache;
private final Logger log;
private final Sensor createTaskSensor;
@@ -49,12 +51,14 @@ class StandbyTaskCreator {
final StreamsConfig applicationConfig,
final StreamsMetricsImpl streamsMetrics,
final StateDirectory stateDirectory,
+ final Time time,
final String threadId,
final LogContext logContext) {
this.topologyMetadata = topologyMetadata;
this.applicationConfig = applicationConfig;
this.streamsMetrics = streamsMetrics;
this.stateDirectory = stateDirectory;
+ this.time = time;
this.log = logContext.logger(getClass());
createTaskSensor = ThreadMetrics.createTaskSensor(threadId,
streamsMetrics);
@@ -85,6 +89,7 @@ class StandbyTaskCreator {
applicationConfig.getBoolean(StreamsConfig.TRANSACTIONAL_STATE_STORES_CONFIG),
getLogContext(taskId),
stateDirectory,
+ time,
topology.storeToChangelogTopic(),
partitions,
upgradeFrom);
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 e1c018aa809..1dee67e509c 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
@@ -257,6 +257,7 @@ public class StateDirectory implements AutoCloseable {
eosEnabled,
logContext,
this,
+ time,
subTopology.storeToChangelogTopic(),
inputPartitions
);
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 bb85da626b0..d731215f9d3 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
@@ -467,6 +467,7 @@ public class StreamThread extends Thread implements
ProcessingThread {
config,
streamsMetrics,
stateDirectory,
+ time,
threadId,
logContext);
diff --git
a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java
b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java
index 3e697138ad4..726a697c139 100644
---
a/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/processor/internals/ProcessorStateManagerTest.java
@@ -131,6 +131,7 @@ public class ProcessorStateManagerTest {
private File checkpointFile;
private OffsetCheckpoint checkpoint;
private StateDirectory stateDirectory;
+ private final MockTime time = new MockTime();
@Mock
private StateStoreMetadata storeMetadata;
@@ -147,7 +148,7 @@ public class ProcessorStateManagerTest {
put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:1234");
put(StreamsConfig.STATE_DIR_CONFIG, baseDir.getPath());
}
- }), new MockTime(), true, true);
+ }), time, true, true);
checkpointFile = new
File(stateDirectory.getOrCreateDirectoryForTask(taskId), CHECKPOINT_FILE_NAME);
checkpoint = new OffsetCheckpoint(checkpointFile);
}
@@ -206,6 +207,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTwoTopicName),
@@ -227,6 +229,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTopicName)
@@ -324,6 +327,50 @@ public class ProcessorStateManagerTest {
verify(store).close();
}
+ @Test
+ public void shouldRefreshTaskDirectoryModificationTimeOnClose() {
+ final ProcessorStateManager stateMgr =
getStateManager(Task.TaskType.ACTIVE);
+ final StateStore store = mock(StateStore.class);
+ when(store.name()).thenReturn(persistentStoreName);
+
+ stateMgr.registerStateStores(singletonList(store), context);
+ stateMgr.registerStore(store, noopStateRestoreCallback, null);
+
+ final File taskDir =
stateDirectory.getOrCreateDirectoryForTask(taskId);
+ final long staleTime = time.milliseconds() - 60_000L;
+ assertTrue(taskDir.setLastModified(staleTime));
+ assertThat(taskDir.lastModified(), is(staleTime));
+
+ stateMgr.close();
+
+ assertThat(taskDir.lastModified(), is(time.milliseconds()));
+ }
+
+ @Test
+ public void
shouldNotRefreshTaskDirectoryModificationTimeWhenClosingStartupTask() {
+ final ProcessorStateManager stateMgr =
ProcessorStateManager.createStartupTaskStateManager(
+ taskId,
+ false,
+ logContext,
+ stateDirectory,
+ time,
+ mkMap(mkEntry(persistentStoreName, persistentStoreTopicName)),
+ emptySet());
+ final StateStore store = mock(StateStore.class);
+ when(store.name()).thenReturn(persistentStoreName);
+
+ stateMgr.registerStateStores(singletonList(store), context);
+ stateMgr.registerStore(store, noopStateRestoreCallback, null);
+
+ final File taskDir =
stateDirectory.getOrCreateDirectoryForTask(taskId);
+ final long staleTime = time.milliseconds() - 60_000L;
+ assertTrue(taskDir.setLastModified(staleTime));
+
+ stateMgr.close();
+
+ assertThat(taskDir.lastModified(), is(staleTime));
+ }
+
@Test
public void shouldRecycleAndReinitializeStore() {
final ProcessorStateManager stateMgr =
getStateManager(Task.TaskType.ACTIVE);
@@ -403,6 +450,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
emptyMap(),
emptySet()
);
@@ -689,6 +737,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
emptyMap(),
emptySet());
@@ -1313,6 +1362,7 @@ public class ProcessorStateManagerTest {
transactionalStateStoresEnabled,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTwoTopicName),
@@ -1330,6 +1380,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTwoTopicName),
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 2c0223fc420..dd37e70ef02 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
@@ -4639,6 +4639,7 @@ public class StreamThreadTest {
config,
streamsMetrics,
stateDirectory,
+ mockTime,
CLIENT_ID,
logContext);
return standbyTaskCreator.createTasks(singletonMap(new TaskId(1, 2),
emptySet()));
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 6bbe268b34f..a6dfe4b29ce 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
@@ -658,6 +658,7 @@ public class StreamThreadStateStoreProviderTest {
false,
logContext,
stateDirectory,
+ new MockTime(),
topology.storeToChangelogTopic(),
partitions);
final RecordCollector recordCollector = new RecordCollectorImpl(
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 597db0fed1b..54020c50de9 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
@@ -502,6 +502,7 @@ public class TopologyTestDriver implements Closeable {
streamsConfig.getBoolean(StreamsConfig.TRANSACTIONAL_STATE_STORES_CONFIG),
logContext,
stateDirectory,
+ mockWallClockTime,
processorTopology.storeToChangelogTopic(),
new HashSet<>(partitionsByInputTopic.values()));
final RecordCollector recordCollector = new RecordCollectorImpl(