bbejeck commented on code in PR #22625:
URL: https://github.com/apache/kafka/pull/22625#discussion_r3443406874
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java:
##########
@@ -336,7 +343,8 @@ void openRocksDB(final DBOptions dbOptions,
new ColumnFamilyDescriptor(OFFSETS_COLUMN_FAMILY_NAME,
offsetsCFOptions())
);
- cfAccessor = new SingleColumnFamilyAccessor(columnFamilies.get(1),
columnFamilies.get(0));
+ dataColumnFamily = columnFamilies.get(0);
Review Comment:
I thought the first column family was for offsets?
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTransactionBuffer.java:
##########
@@ -0,0 +1,323 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.streams.state.internals;
+
+import org.apache.kafka.common.utils.Bytes;
+import org.apache.kafka.streams.KeyValue;
+import org.apache.kafka.streams.errors.ProcessorStateException;
+
+import org.rocksdb.ColumnFamilyHandle;
+import org.rocksdb.RocksDB;
+import org.rocksdb.RocksDBException;
+import org.rocksdb.RocksIterator;
+import org.rocksdb.WriteBatch;
+import org.rocksdb.WriteOptions;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.NavigableMap;
+import java.util.NoSuchElementException;
+import java.util.Optional;
+import java.util.TreeMap;
+
+/**
+ * A {@link TransactionBuffer} implementation for RocksDB-backed stores.
+ * Uses a {@link WriteBatch} (without index) to accumulate writes for atomic
commit.
+ * Reads are handled entirely by the shared {@code ConcurrentSkipListMap} in
+ * {@link AbstractTransactionBuffer}, so a {@code WriteBatchWithIndex} is not
needed.
+ * <p>
+ * Range deletions ({@link #stageDeleteRange}) are only supported by
RocksDB-backed stores
+ * and are owned entirely by this class; {@link AbstractTransactionBuffer}
carries no
+ * tombstone state.
+ */
+class RocksDBTransactionBuffer extends AbstractTransactionBuffer<Bytes> {
+
+ private final RocksDB db;
+ private final ColumnFamilyHandle cfHandle;
+ private final WriteOptions wOptions;
+ private final String storeName;
+ private WriteBatch writeBatch;
+ private volatile NavigableMap<Bytes, List<Bytes>> rangeTombstones =
Collections.emptyNavigableMap();
+
+ RocksDBTransactionBuffer(final RocksDB db,
+ final ColumnFamilyHandle cfHandle,
+ final WriteOptions wOptions,
+ final String storeName) {
+ this.db = db;
+ this.cfHandle = cfHandle;
+ this.wOptions = wOptions;
+ this.storeName = storeName;
+ this.writeBatch = new WriteBatch();
+ }
+
+ @Override
+ int estimateKeySize(final Bytes key) {
+ return key.get().length;
+ }
+
+ @Override
+ void stageToBackend(final Bytes key, final byte[] value) {
+ stage(cfHandle, key, value);
+ }
+
+ /**
+ * Stages a write for an explicit column family. Updates the shared read
buffer
+ * ({@code pendingWrites}) and appends the write to the shared {@link
WriteBatch}
+ * under {@code cf}, so every staged CF is committed atomically on {@link
#commit()}.
+ */
+ void stage(final ColumnFamilyHandle cf, final Bytes key, final byte[]
value) {
+ pendingWrites.put(key, Optional.ofNullable(value));
+ pendingWritesBytes += estimateKeySize(key) + (value != null ?
value.length : 0);
+ try {
+ if (value != null) {
+ writeBatch.put(cf, key.get(), value);
+ } else {
+ writeBatch.delete(cf, key.get());
+ }
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging write in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ /**
+ * Stages a range deletion for an explicit column family. Updates the
shared
+ * {@code pendingWrites} and {@code rangeTombstones} so iterators opened
before
+ * commit hide the deleted range, and appends the range delete to the
shared
+ * {@link WriteBatch} under {@code cf}.
+ */
+ void stageDeleteRange(final ColumnFamilyHandle cf, final Bytes from, final
Bytes to) {
+ pendingWrites.subMap(from, true, to, false).clear();
+ final TreeMap<Bytes, List<Bytes>> copy = new
TreeMap<>(rangeTombstones);
+ copy.computeIfAbsent(from, k -> new ArrayList<>()).add(to);
+ rangeTombstones = copy;
+ pendingWritesBytes += estimateKeySize(from) + estimateKeySize(to);
+ try {
+ writeBatch.deleteRange(cf, from.get(), to.get());
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging delete range in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ @Override
+ public Optional<byte[]> get(final Bytes key) {
+ final Optional<byte[]> staged = pendingWrites.get(key);
+ if (staged != null) {
Review Comment:
same here
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTransactionBuffer.java:
##########
@@ -0,0 +1,323 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.streams.state.internals;
+
+import org.apache.kafka.common.utils.Bytes;
+import org.apache.kafka.streams.KeyValue;
+import org.apache.kafka.streams.errors.ProcessorStateException;
+
+import org.rocksdb.ColumnFamilyHandle;
+import org.rocksdb.RocksDB;
+import org.rocksdb.RocksDBException;
+import org.rocksdb.RocksIterator;
+import org.rocksdb.WriteBatch;
+import org.rocksdb.WriteOptions;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.NavigableMap;
+import java.util.NoSuchElementException;
+import java.util.Optional;
+import java.util.TreeMap;
+
+/**
+ * A {@link TransactionBuffer} implementation for RocksDB-backed stores.
+ * Uses a {@link WriteBatch} (without index) to accumulate writes for atomic
commit.
+ * Reads are handled entirely by the shared {@code ConcurrentSkipListMap} in
+ * {@link AbstractTransactionBuffer}, so a {@code WriteBatchWithIndex} is not
needed.
+ * <p>
+ * Range deletions ({@link #stageDeleteRange}) are only supported by
RocksDB-backed stores
+ * and are owned entirely by this class; {@link AbstractTransactionBuffer}
carries no
+ * tombstone state.
+ */
+class RocksDBTransactionBuffer extends AbstractTransactionBuffer<Bytes> {
+
+ private final RocksDB db;
+ private final ColumnFamilyHandle cfHandle;
+ private final WriteOptions wOptions;
+ private final String storeName;
+ private WriteBatch writeBatch;
+ private volatile NavigableMap<Bytes, List<Bytes>> rangeTombstones =
Collections.emptyNavigableMap();
+
+ RocksDBTransactionBuffer(final RocksDB db,
+ final ColumnFamilyHandle cfHandle,
+ final WriteOptions wOptions,
+ final String storeName) {
+ this.db = db;
+ this.cfHandle = cfHandle;
+ this.wOptions = wOptions;
+ this.storeName = storeName;
+ this.writeBatch = new WriteBatch();
+ }
+
+ @Override
+ int estimateKeySize(final Bytes key) {
+ return key.get().length;
+ }
+
+ @Override
+ void stageToBackend(final Bytes key, final byte[] value) {
+ stage(cfHandle, key, value);
+ }
+
+ /**
+ * Stages a write for an explicit column family. Updates the shared read
buffer
+ * ({@code pendingWrites}) and appends the write to the shared {@link
WriteBatch}
+ * under {@code cf}, so every staged CF is committed atomically on {@link
#commit()}.
+ */
+ void stage(final ColumnFamilyHandle cf, final Bytes key, final byte[]
value) {
+ pendingWrites.put(key, Optional.ofNullable(value));
+ pendingWritesBytes += estimateKeySize(key) + (value != null ?
value.length : 0);
+ try {
+ if (value != null) {
+ writeBatch.put(cf, key.get(), value);
+ } else {
+ writeBatch.delete(cf, key.get());
+ }
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging write in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ /**
+ * Stages a range deletion for an explicit column family. Updates the
shared
+ * {@code pendingWrites} and {@code rangeTombstones} so iterators opened
before
+ * commit hide the deleted range, and appends the range delete to the
shared
+ * {@link WriteBatch} under {@code cf}.
+ */
+ void stageDeleteRange(final ColumnFamilyHandle cf, final Bytes from, final
Bytes to) {
+ pendingWrites.subMap(from, true, to, false).clear();
+ final TreeMap<Bytes, List<Bytes>> copy = new
TreeMap<>(rangeTombstones);
+ copy.computeIfAbsent(from, k -> new ArrayList<>()).add(to);
+ rangeTombstones = copy;
+ pendingWritesBytes += estimateKeySize(from) + estimateKeySize(to);
+ try {
+ writeBatch.deleteRange(cf, from.get(), to.get());
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging delete range in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ @Override
+ public Optional<byte[]> get(final Bytes key) {
+ final Optional<byte[]> staged = pendingWrites.get(key);
+ if (staged != null) {
+ return staged;
+ }
+ if (isCoveredByRangeTombstone(key, rangeTombstones)) {
+ return Optional.empty();
+ }
+ return null;
+ }
+
+ @Override
+ public boolean isEmpty() {
+ return super.isEmpty() && rangeTombstones.isEmpty();
+ }
+
+ ManagedKeyValueIterator<Bytes, byte[]> all(final ColumnFamilyHandle cf,
final boolean forward) {
+ if (Thread.currentThread() == ownerThread) {
+ final ManagedKeyValueIterator<Bytes, byte[]> baseIter =
newBaseIterator(cf, null, null, forward, true);
+ return new StagedMergeIterator<>(pendingWrites, baseIter, forward);
+ }
+ return snapshotScan(cf, null, null, forward, true);
+ }
+
+ ManagedKeyValueIterator<Bytes, byte[]> range(final ColumnFamilyHandle cf,
+ final Bytes from, final Bytes
to,
+ final boolean forward, final
boolean toInclusive) {
+ if (Thread.currentThread() == ownerThread) {
+ final ManagedKeyValueIterator<Bytes, byte[]> baseIter =
newBaseIterator(cf, from, to, forward, toInclusive);
+ final NavigableMap<Bytes, Optional<byte[]>> stagingView =
boundStaging(from, to, toInclusive);
+ return new StagedMergeIterator<>(stagingView, baseIter, forward);
+ }
+ return snapshotScan(cf, from, to, forward, toInclusive);
+ }
+
+ private ManagedKeyValueIterator<Bytes, byte[]> snapshotScan(final
ColumnFamilyHandle cf,
+ final Bytes
from, final Bytes to,
+ final boolean
forward, final boolean toInclusive) {
+ snapshotLock.readLock().lock();
+ try {
+ final NavigableMap<Bytes, Optional<byte[]>> stagingSnapshot =
+ new TreeMap<>(boundStaging(from, to, toInclusive));
+ final ManagedKeyValueIterator<Bytes, byte[]> baseIter =
newBaseIterator(cf, from, to, forward, toInclusive);
+ return new StagedMergeIterator<>(stagingSnapshot, baseIter,
forward);
+ } finally {
+ snapshotLock.readLock().unlock();
+ }
+ }
+
+ @Override
+ ManagedKeyValueIterator<Bytes, byte[]> newBaseIterator(final Bytes from,
final Bytes to) {
+ return newBaseIterator(cfHandle, from, to, true, true);
+ }
+
+ @Override
+ ManagedKeyValueIterator<Bytes, byte[]> newBaseIterator(final Bytes from,
final Bytes to,
+ final boolean
forward, final boolean toInclusive) {
+ return newBaseIterator(cfHandle, from, to, forward, toInclusive);
+ }
+
+ private ManagedKeyValueIterator<Bytes, byte[]> newBaseIterator(final
ColumnFamilyHandle cf,
+ final Bytes
from, final Bytes to,
+ final
boolean forward, final boolean toInclusive) {
+ final RocksIterator rocksIterator = db.newIterator(cf);
+ final ManagedKeyValueIterator<Bytes, byte[]> iter;
+ if (from != null && to != null) {
+ iter = new RocksDBRangeIterator(storeName, rocksIterator, from,
to, forward, toInclusive);
+ } else if (from != null && forward) {
+ rocksIterator.seek(from.get());
+ iter = new RocksDbIterator(storeName, rocksIterator, true);
+ } else if (!forward) {
+ if (to != null) {
+ rocksIterator.seekForPrev(to.get());
+ } else {
+ rocksIterator.seekToLast();
+ }
+ iter = new RocksDbIterator(storeName, rocksIterator, false);
+ } else {
+ rocksIterator.seekToFirst();
+ iter = new RocksDbIterator(storeName, rocksIterator, true);
+ }
+ // RocksDbIterator requires onClose to be set before close() is called.
+ // Since this iterator is used internally by StagedMergeIterator (not
+ // tracked by RocksDBStore's open-iterator set), use a no-op callback.
+ iter.onClose(() -> { });
+ if (rangeTombstones.isEmpty()) {
+ return iter;
+ }
+ return new RangeTombstoneFilterIterator(iter, rangeTombstones);
+ }
+
+ @Override
+ void flushToBase() {
+ try {
+ db.write(wOptions, writeBatch);
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error committing transaction
buffer for store " + storeName, e);
+ }
+ writeBatch.close();
+ writeBatch = new WriteBatch();
+ rangeTombstones = Collections.emptyNavigableMap();
+ }
+
+ @Override
+ void discardPendingBatch() {
+ writeBatch.clear();
+ rangeTombstones = Collections.emptyNavigableMap();
+ }
+
+ @Override
+ public long approximateNumUncommittedBytes() {
+ return super.approximateNumUncommittedBytes() +
writeBatch.getDataSize();
Review Comment:
This seems to double count - once with `pendingWritesBytes` line 86 and the
serialized batch which contains the same bytes - so it seems that would reduce
the max uncommited bytes config in half.
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBStore.java:
##########
@@ -1034,6 +1050,115 @@ public void close() {
}
+ static class TransactionalDBAccessor implements DBAccessor {
+
+ private final DBAccessor underlying;
+ private final RocksDBTransactionBuffer buffer;
+ private final ColumnFamilyHandle cfHandle;
+
+ TransactionalDBAccessor(final DBAccessor underlying,
+ final RocksDB db,
+ final ColumnFamilyHandle cfHandle,
+ final WriteOptions wOptions,
+ final String storeName) {
+ this.underlying = underlying;
+ this.cfHandle = cfHandle;
+ this.buffer = new RocksDBTransactionBuffer(db, cfHandle, wOptions,
storeName);
+ }
+
+ @Override
+ public byte[] get(final ColumnFamilyHandle columnFamily, final byte[]
key) throws RocksDBException {
+ if (columnFamily.equals(cfHandle)) {
+ final java.util.Optional<byte[]> staged =
buffer.get(Bytes.wrap(key));
+ if (staged != null) {
+ return staged.orElse(null);
+ }
+ }
+ return underlying.get(columnFamily, key);
+ }
+
+ @Override
+ public byte[] get(final ColumnFamilyHandle columnFamily, final
ReadOptions readOptions, final byte[] key) throws RocksDBException {
+ if (columnFamily.equals(cfHandle)) {
+ final java.util.Optional<byte[]> staged =
buffer.get(Bytes.wrap(key));
Review Comment:
if `buffer.get` returns an `Optional` why the `null` check? I say this
becuase below we use `orElse`
##########
streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDBTransactionBuffer.java:
##########
@@ -0,0 +1,323 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.streams.state.internals;
+
+import org.apache.kafka.common.utils.Bytes;
+import org.apache.kafka.streams.KeyValue;
+import org.apache.kafka.streams.errors.ProcessorStateException;
+
+import org.rocksdb.ColumnFamilyHandle;
+import org.rocksdb.RocksDB;
+import org.rocksdb.RocksDBException;
+import org.rocksdb.RocksIterator;
+import org.rocksdb.WriteBatch;
+import org.rocksdb.WriteOptions;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.NavigableMap;
+import java.util.NoSuchElementException;
+import java.util.Optional;
+import java.util.TreeMap;
+
+/**
+ * A {@link TransactionBuffer} implementation for RocksDB-backed stores.
+ * Uses a {@link WriteBatch} (without index) to accumulate writes for atomic
commit.
+ * Reads are handled entirely by the shared {@code ConcurrentSkipListMap} in
+ * {@link AbstractTransactionBuffer}, so a {@code WriteBatchWithIndex} is not
needed.
+ * <p>
+ * Range deletions ({@link #stageDeleteRange}) are only supported by
RocksDB-backed stores
+ * and are owned entirely by this class; {@link AbstractTransactionBuffer}
carries no
+ * tombstone state.
+ */
+class RocksDBTransactionBuffer extends AbstractTransactionBuffer<Bytes> {
+
+ private final RocksDB db;
+ private final ColumnFamilyHandle cfHandle;
+ private final WriteOptions wOptions;
+ private final String storeName;
+ private WriteBatch writeBatch;
+ private volatile NavigableMap<Bytes, List<Bytes>> rangeTombstones =
Collections.emptyNavigableMap();
+
+ RocksDBTransactionBuffer(final RocksDB db,
+ final ColumnFamilyHandle cfHandle,
+ final WriteOptions wOptions,
+ final String storeName) {
+ this.db = db;
+ this.cfHandle = cfHandle;
+ this.wOptions = wOptions;
+ this.storeName = storeName;
+ this.writeBatch = new WriteBatch();
+ }
+
+ @Override
+ int estimateKeySize(final Bytes key) {
+ return key.get().length;
+ }
+
+ @Override
+ void stageToBackend(final Bytes key, final byte[] value) {
+ stage(cfHandle, key, value);
+ }
+
+ /**
+ * Stages a write for an explicit column family. Updates the shared read
buffer
+ * ({@code pendingWrites}) and appends the write to the shared {@link
WriteBatch}
+ * under {@code cf}, so every staged CF is committed atomically on {@link
#commit()}.
+ */
+ void stage(final ColumnFamilyHandle cf, final Bytes key, final byte[]
value) {
+ pendingWrites.put(key, Optional.ofNullable(value));
+ pendingWritesBytes += estimateKeySize(key) + (value != null ?
value.length : 0);
+ try {
+ if (value != null) {
+ writeBatch.put(cf, key.get(), value);
+ } else {
+ writeBatch.delete(cf, key.get());
+ }
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging write in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ /**
+ * Stages a range deletion for an explicit column family. Updates the
shared
+ * {@code pendingWrites} and {@code rangeTombstones} so iterators opened
before
+ * commit hide the deleted range, and appends the range delete to the
shared
+ * {@link WriteBatch} under {@code cf}.
+ */
+ void stageDeleteRange(final ColumnFamilyHandle cf, final Bytes from, final
Bytes to) {
+ pendingWrites.subMap(from, true, to, false).clear();
+ final TreeMap<Bytes, List<Bytes>> copy = new
TreeMap<>(rangeTombstones);
+ copy.computeIfAbsent(from, k -> new ArrayList<>()).add(to);
+ rangeTombstones = copy;
+ pendingWritesBytes += estimateKeySize(from) + estimateKeySize(to);
+ try {
+ writeBatch.deleteRange(cf, from.get(), to.get());
+ } catch (final RocksDBException e) {
+ throw new ProcessorStateException("Error staging delete range in
transaction buffer for store " + storeName, e);
+ }
+ }
+
+ @Override
+ public Optional<byte[]> get(final Bytes key) {
+ final Optional<byte[]> staged = pendingWrites.get(key);
+ if (staged != null) {
+ return staged;
+ }
+ if (isCoveredByRangeTombstone(key, rangeTombstones)) {
+ return Optional.empty();
+ }
+ return null;
+ }
+
+ @Override
+ public boolean isEmpty() {
+ return super.isEmpty() && rangeTombstones.isEmpty();
+ }
+
+ ManagedKeyValueIterator<Bytes, byte[]> all(final ColumnFamilyHandle cf,
final boolean forward) {
+ if (Thread.currentThread() == ownerThread) {
+ final ManagedKeyValueIterator<Bytes, byte[]> baseIter =
newBaseIterator(cf, null, null, forward, true);
+ return new StagedMergeIterator<>(pendingWrites, baseIter, forward);
+ }
+ return snapshotScan(cf, null, null, forward, true);
+ }
+
+ ManagedKeyValueIterator<Bytes, byte[]> range(final ColumnFamilyHandle cf,
+ final Bytes from, final Bytes
to,
+ final boolean forward, final
boolean toInclusive) {
+ if (Thread.currentThread() == ownerThread) {
+ final ManagedKeyValueIterator<Bytes, byte[]> baseIter =
newBaseIterator(cf, from, to, forward, toInclusive);
+ final NavigableMap<Bytes, Optional<byte[]>> stagingView =
boundStaging(from, to, toInclusive);
+ return new StagedMergeIterator<>(stagingView, baseIter, forward);
Review Comment:
if `boundStaging` throws `baseIter` leaks, switch the creation order?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]