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 e409145cb39 KAFKA-20753: Fix transactional status marker writes
(#22716)
e409145cb39 is described below
commit e409145cb396b61337922e65a56f98f434d641a7
Author: Nick Telford <[email protected]>
AuthorDate: Wed Jul 1 16:38:50 2026 +0100
KAFKA-20753: Fix transactional status marker writes (#22716)
When `enable.transactional.statestores=true`, the KIP-1035 open/closed
status marker written to the RocksDB offsets column family on `open()`
and `close()` was being routed through `TransactionalDBAccessor.put()`,
which stages writes in the in-memory `RocksDBTransactionBuffer`. The
buffer is discarded without committing on store close, so the CLOSED
marker was silently lost. The on-disk status remained permanently at
OPEN, causing every subsequent `open()` under EOS
(`ignoreInvalidState=false`) to throw `ProcessorStateException("Invalid
state during store open...")`, which escalated to
`TaskCorruptedException`. Because the wipe-and-restore cycle would
reproduce the same condition, this created an unrecoverable loop
(KAFKA-20753).
The fix suppresses status marker writes in
`AbstractColumnFamilyAccessor` when the store is transactional
(`isTransactional=true`), while preserving the read-and-check on
`open()`. This means a stale OPEN marker left by a prior
non-transactional crash — e.g. after toggling
`enable.transactional.statestores` — is still detected and triggers the
correct EOS wipe. Writing the marker is unnecessary for transactional
stores because KIP-892's atomic `WriteBatch` commit ensures uncommitted
data is never persisted to disk; the original safety guarantee the
marker provided no longer applies.
The `isTransactional` flag is set from configs in
`RocksDBStore.openDB()` before `openRocksDB()` is called, so all
`ColumnFamilyAccessor` subclasses (`SingleColumnFamilyAccessor`,
`DualColumnFamilyAccessor`) pick it up at construction time. Three new
tests are added to `AbstractColumnFamilyAccessorTest` covering: no OPEN
marker written on open, no CLOSED marker written on close, and detection
of a stale OPEN from a prior non-transactional crash.
Reviewers: Bill Bejeck <[email protected]>
---
.../internals/AbstractColumnFamilyAccessor.java | 14 ++++++---
.../state/internals/DualColumnFamilyAccessor.java | 2 +-
.../streams/state/internals/RocksDBStore.java | 10 ++++---
.../AbstractColumnFamilyAccessorTest.java | 35 +++++++++++++++++++++-
.../internals/DualColumnFamilyAccessorTest.java | 13 ++++++++
.../state/internals/RocksDBStoreCloseLeakTest.java | 8 ++---
6 files changed, 68 insertions(+), 14 deletions(-)
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessor.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessor.java
index c24b689885e..aebcecd75a7 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessor.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessor.java
@@ -46,10 +46,15 @@ abstract class AbstractColumnFamilyAccessor implements
RocksDBStore.ColumnFamily
private final byte[] openState = longSerde.serializer().serialize(null,
1L);
private final byte[] closedState = longSerde.serializer().serialize(null,
0L);
private final AtomicBoolean storeOpen;
+ // Status marker writes are skipped when true; the check on open still
runs.
+ private final boolean isTransactional;
- AbstractColumnFamilyAccessor(final ColumnFamilyHandle
offsetColumnFamilyHandle, final AtomicBoolean storeOpen) {
+ AbstractColumnFamilyAccessor(final ColumnFamilyHandle
offsetColumnFamilyHandle,
+ final AtomicBoolean storeOpen,
+ final boolean isTransactional) {
this.offsetColumnFamilyHandle = offsetColumnFamilyHandle;
this.storeOpen = storeOpen;
+ this.isTransactional = isTransactional;
}
@Override
@@ -81,8 +86,9 @@ abstract class AbstractColumnFamilyAccessor implements
RocksDBStore.ColumnFamily
public final Position open(final RocksDBStore.DBAccessor accessor, final
boolean ignoreInvalidState) throws RocksDBException {
final byte[] valueBytes = accessor.get(offsetColumnFamilyHandle,
statusKey);
if (ignoreInvalidState || (valueBytes == null ||
Arrays.equals(valueBytes, closedState))) {
- // If the status key is not present, we initialize it to "OPEN"
- accessor.put(offsetColumnFamilyHandle, statusKey, openState);
+ if (!isTransactional) {
+ accessor.put(offsetColumnFamilyHandle, statusKey, openState);
+ }
storeOpen.set(true);
final byte[] positionBytes =
accessor.get(offsetColumnFamilyHandle, positionKey);
if (positionBytes != null) {
@@ -103,7 +109,7 @@ abstract class AbstractColumnFamilyAccessor implements
RocksDBStore.ColumnFamily
// failed state (e.g. during an EOSv2 fencing cascade); the handle
close must
// still happen, otherwise the native ColumnFamilyHandle leaks every
cycle.
try {
- if (storeOpen.compareAndSet(true, false)) {
+ if (storeOpen.compareAndSet(true, false) && !isTransactional) {
accessor.put(offsetColumnFamilyHandle, statusKey, closedState);
}
} finally {
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java
index 4b2f7ba5bf4..c41aaf4a357 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessor.java
@@ -70,7 +70,7 @@ class DualColumnFamilyAccessor extends
AbstractColumnFamilyAccessor {
final Function<byte[], byte[]> valueConverter,
final RocksDBStore store,
final AtomicBoolean storeOpen) {
- super(offsetColumnFamily, storeOpen);
+ super(offsetColumnFamily, storeOpen, store.isTransactional);
this.oldColumnFamily = oldColumnFamily;
this.newColumnFamily = newColumnFamily;
this.valueConverter = valueConverter;
diff --git
a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java
b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java
index 417c36e8547..4edd65ae8cd 100644
---
a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java
+++
b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java
@@ -123,6 +123,8 @@ public class RocksDBStore implements KeyValueStore<Bytes,
byte[]>, BatchWritingS
DBAccessor dbAccessor;
ColumnFamilyAccessor cfAccessor;
protected final AtomicBoolean open = new AtomicBoolean(false);
+ // package-private: read by DualColumnFamilyAccessor
+ boolean isTransactional;
// the following option objects will be created in openDB and closed in
the close() method
private RocksDBGenericOptionsToDbOptionsColumnFamilyOptionsAdapter
userSpecifiedOptions;
@@ -256,6 +258,8 @@ public class RocksDBStore implements KeyValueStore<Bytes,
byte[]>, BatchWritingS
throw new ProcessorStateException(fatal);
}
+ isTransactional = StreamsConfig.InternalConfig.getBoolean(configs,
StreamsConfig.TRANSACTIONAL_STATE_STORES_CONFIG, false);
+
// Setup statistics before the database is opened, otherwise the
statistics are not updated
// with the measurements from Rocks DB
setupStatistics(configs, dbOptions);
@@ -285,9 +289,7 @@ public class RocksDBStore implements KeyValueStore<Bytes,
byte[]>, BatchWritingS
throw e;
}
- final boolean transactional = StreamsConfig.InternalConfig.getBoolean(
- configs, StreamsConfig.TRANSACTIONAL_STATE_STORES_CONFIG, false);
- if (transactional) {
+ if (isTransactional) {
dbAccessor = new TransactionalDBAccessor(dbAccessor, db,
cfAccessor.dataColumnFamily(), cfAccessor.offsetsColumnFamily(), wOptions,
name);
}
@@ -1408,7 +1410,7 @@ public class RocksDBStore implements KeyValueStore<Bytes,
byte[]>, BatchWritingS
private final ColumnFamilyHandle columnFamily;
SingleColumnFamilyAccessor(final ColumnFamilyHandle
offsetsColumnFamily, final ColumnFamilyHandle columnFamily) {
- super(offsetsColumnFamily, open);
+ super(offsetsColumnFamily, open, isTransactional);
this.columnFamily = columnFamily;
}
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessorTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessorTest.java
index 848030dc570..198f775147d 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessorTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/AbstractColumnFamilyAccessorTest.java
@@ -61,6 +61,8 @@ abstract class AbstractColumnFamilyAccessorTest {
protected AbstractColumnFamilyAccessor accessor;
abstract AbstractColumnFamilyAccessor createColumnFamilyAccessor();
+
+ abstract AbstractColumnFamilyAccessor
createTransactionalColumnFamilyAccessor();
private final LongSerializer offsetSerializer = new LongSerializer();
private final StringSerializer keySerializer = new StringSerializer();
private final byte[] openValue = toBytes(1L);
@@ -170,10 +172,41 @@ abstract class AbstractColumnFamilyAccessorTest {
verify(dbAccessor, never()).put(any(), any(), any());
}
+ @Test
+ public void shouldNotWriteOpenMarkerForTransactionalStore() throws
RocksDBException {
+ dbAccessor = new InMemoryRocksDBAccessor(mock(RocksDB.class));
+ final AbstractColumnFamilyAccessor txnAccessor =
createTransactionalColumnFamilyAccessor();
+ txnAccessor.open(dbAccessor, false);
+ assertNull(dbAccessor.get(offsetsCF, toBytes("status")));
+ }
+
+ @Test
+ public void shouldNotWriteClosedMarkerForTransactionalStore() throws
RocksDBException {
+ dbAccessor = new InMemoryRocksDBAccessor(mock(RocksDB.class));
+ final AbstractColumnFamilyAccessor txnAccessor =
createTransactionalColumnFamilyAccessor();
+ txnAccessor.open(dbAccessor, false);
+ txnAccessor.close(dbAccessor);
+ assertNull(dbAccessor.get(offsetsCF, toBytes("status")));
+ }
+
+ @Test
+ public void shouldDetectStaleOpenMarkerFromPriorNonTransactionalCrash()
throws RocksDBException {
+ dbAccessor = new InMemoryRocksDBAccessor(mock(RocksDB.class));
+ // Simulate a prior non-transactional crash: OPEN marker left on disk
+ dbAccessor.put(offsetsCF, toBytes("status"), openValue);
+
+ final AbstractColumnFamilyAccessor txnAccessor =
createTransactionalColumnFamilyAccessor();
+ final ProcessorStateException thrown = assertThrowsExactly(
+ ProcessorStateException.class, () ->
txnAccessor.open(dbAccessor, false));
+ assertEquals("Invalid state during store open. Expected state to be
either empty or closed", thrown.getMessage());
+ // Marker is unchanged — transactional accessor never writes status
+ assertArrayEquals(openValue, dbAccessor.get(offsetsCF,
toBytes("status")));
+ }
+
private byte[] toBytes(final String s) {
return keySerializer.serialize("", s);
}
-
+
private byte[] toBytes(final long l) {
return offsetSerializer.serialize("", l);
}
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessorTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessorTest.java
index ba8be285ac3..c932c012abd 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessorTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/DualColumnFamilyAccessorTest.java
@@ -91,6 +91,19 @@ public class DualColumnFamilyAccessorTest extends
AbstractColumnFamilyAccessorTe
return new DualColumnFamilyAccessor(offsetsCF, oldCF, newCF,
valueConverter, store, storeOpen);
}
+ @Override
+ AbstractColumnFamilyAccessor createTransactionalColumnFamilyAccessor() {
+ final RocksDBStore store = mock(RocksDBStore.class);
+ store.position = Position.emptyPosition();
+ store.context = mock(StateStoreContext.class);
+ store.isTransactional = true;
+ lenient().when(store.name()).thenReturn(STORE_NAME);
+ if (valueConverter == null) {
+ valueConverter = oldValue -> oldValue == null ? null
+ : ByteBuffer.allocate(oldValue.length +
10).put("converted:".getBytes()).put(oldValue).array();
+ }
+ return new DualColumnFamilyAccessor(offsetsCF, oldCF, newCF,
valueConverter, store, storeOpen);
+ }
@Test
public void shouldPutValueToNewColumnFamilyAndDeleteFromOld() throws
RocksDBException {
diff --git
a/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBStoreCloseLeakTest.java
b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBStoreCloseLeakTest.java
index 2e3e2242a87..b709c9c55a1 100644
---
a/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBStoreCloseLeakTest.java
+++
b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDBStoreCloseLeakTest.java
@@ -45,10 +45,10 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
* Constructing a {@code ColumnFamilyOptions} on the JNI side
auto-allocates a default
* {@code BlockBasedTableFactory} and its {@code LRUCache}; the cache
itself has no Java
* handle, so we assert on the options' own handle and rely on the
destructor contract.</li>
- * <li>KIP-1035 close-path — {@link AbstractColumnFamilyAccessor#close} writes
a closed-state
- * marker to the offsets CF. If that {@code put} throws (e.g. during an
EOSv2 fencing
- * cascade or unclean shutdown), the column family handles must still be
released. We
- * simulate the throw via {@link ThrowingOnOffsetsPutDBAccessor} and
observe via
+ * <li>KIP-1035 close-path — for non-transactional stores, {@link
AbstractColumnFamilyAccessor#close}
+ * writes a closed-state marker to the offsets CF. If that {@code put}
throws (e.g. during
+ * an EOSv2 fencing cascade or unclean shutdown), the column family
handles must still be
+ * released. We simulate the throw via {@link
ThrowingOnOffsetsPutDBAccessor} and observe via
* {@code isOwningHandle()} because {@code RocksDBStore.close()} swallows
the resulting
* {@code RocksDBException}.</li>
* </ol>