This is an automated email from the ASF dual-hosted git repository.
bbejeck pushed a commit to branch 4.3
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/4.3 by this push:
new b27b48dbf60 KAFKA-20805: Apply fix for refreshing task directory mtime
on state manager close on 4.3 (#22844)
b27b48dbf60 is described below
commit b27b48dbf608460f5ec2f84333b2cbbfaaa09d98
Author: Bill Bejeck <[email protected]>
AuthorDate: Wed Jul 15 20:10:08 2026 -0400
KAFKA-20805: Apply fix for refreshing task directory mtime on state manager
close on 4.3 (#22844)
Port of the trunk PR #22837 and the original author is @nicktelford
Reviewers: 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 | 52 +++++++++++++++++++++-
.../processor/internals/StreamThreadTest.java | 1 +
.../StreamThreadStateStoreProviderTest.java | 1 +
.../apache/kafka/streams/TopologyTestDriver.java | 1 +
9 files changed, 82 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 1d5b7fbf7ed..2a722657baf 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
@@ -151,6 +151,7 @@ class ActiveTaskCreator {
eosEnabled(applicationConfig),
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 7b3483f8e68..afe393b115e 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
@@ -20,6 +20,7 @@ import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.utils.FixedOrderMap;
import org.apache.kafka.common.utils.LogContext;
+import org.apache.kafka.common.utils.Time;
import org.apache.kafka.streams.errors.ProcessorStateException;
import org.apache.kafka.streams.errors.StreamsException;
import org.apache.kafka.streams.errors.TaskCorruptedException;
@@ -184,6 +185,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;
@@ -205,6 +209,7 @@ public class ProcessorStateManager implements StateManager {
final boolean eosEnabled,
final LogContext logContext,
final StateDirectory stateDirectory,
+ final Time time,
final Map<String, String>
storeToChangelogTopic,
final Collection<TopicPartition>
sourcePartitions,
final UpgradeFromValues upgradeFrom) throws
ProcessorStateException {
@@ -216,6 +221,7 @@ public class ProcessorStateManager implements StateManager {
this.eosEnabled = eosEnabled;
this.sourcePartitions = sourcePartitions;
this.upgradeFrom = upgradeFrom;
+ this.time = time;
this.baseDir = stateDirectory.getOrCreateDirectoryForTask(taskId);
this.stateDirectory = stateDirectory;
@@ -233,9 +239,10 @@ public class ProcessorStateManager implements StateManager
{
final boolean eosEnabled,
final LogContext logContext,
final StateDirectory stateDirectory,
+ final Time time,
final Map<String, String>
storeToChangelogTopic,
final Collection<TopicPartition>
sourcePartitions) throws ProcessorStateException {
- this(taskId, taskType, eosEnabled, logContext, stateDirectory,
storeToChangelogTopic, sourcePartitions, null);
+ this(taskId, taskType, eosEnabled, logContext, stateDirectory, time,
storeToChangelogTopic, sourcePartitions, null);
}
/**
@@ -247,9 +254,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,
logContext, stateDirectory, storeToChangelogTopic, sourcePartitions);
+ final ProcessorStateManager stateManager =
+ new ProcessorStateManager(taskId, TaskType.STANDBY, eosEnabled,
logContext, stateDirectory, time, storeToChangelogTopic, sourcePartitions);
+ stateManager.startupTask = true;
+ return stateManager;
}
void registerStateStores(final List<StateStore> allStores, final
InternalProcessorContext<?, ?> processorContext) {
@@ -646,6 +657,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.debug("{}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 139efbd63de..c221f9fa708 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
@@ -19,6 +19,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.LogContext;
+import org.apache.kafka.common.utils.Time;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.internals.UpgradeFromValues;
import org.apache.kafka.streams.processor.TaskId;
@@ -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);
@@ -84,6 +88,7 @@ class StandbyTaskCreator {
eosEnabled(applicationConfig),
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 c8d68e7c23e..0a4b925e299 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
@@ -255,6 +255,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 833d42aeae4..9ec72d715bd 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
@@ -462,6 +462,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 f2a96bffa3e..8953cce56b8 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);
}
@@ -205,6 +206,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTwoTopicName),
@@ -225,6 +227,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
mkMap(
mkEntry(persistentStoreName, persistentStoreTopicName),
mkEntry(persistentStoreTwoName, persistentStoreTopicName)
@@ -322,6 +325,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);
@@ -400,6 +447,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
emptyMap(),
emptySet()
);
@@ -685,6 +733,7 @@ public class ProcessorStateManagerTest {
false,
logContext,
stateDirectory,
+ time,
emptyMap(),
emptySet());
@@ -1268,6 +1317,7 @@ public class ProcessorStateManagerTest {
eosEnabled,
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 21e06b19cc1..3900c3c97e6 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
@@ -4120,6 +4120,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 59cd8ff8ef1..5dd1cd98374 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
@@ -657,6 +657,7 @@ public class StreamThreadStateStoreProviderTest {
StreamsConfigUtils.eosEnabled(streamsConfig),
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 0d7568b0819..88fb740eef9 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
@@ -488,6 +488,7 @@ public class TopologyTestDriver implements Closeable {
StreamsConfig.EXACTLY_ONCE_V2.equals(streamsConfig.getString(StreamsConfig.PROCESSING_GUARANTEE_CONFIG)),
logContext,
stateDirectory,
+ mockWallClockTime,
processorTopology.storeToChangelogTopic(),
new HashSet<>(partitionsByInputTopic.values()));
final RecordCollector recordCollector = new RecordCollectorImpl(