From 60095b73646e78e16f1f9cebfcc0dce081183b8e Mon Sep 17 00:00:00 2001 From: "Charles P. Wright" Date: Thu, 27 Aug 2026 07:47:19 -0400 Subject: [PATCH 1/3] fix: DH-23494: Close row sets and contexts in TimeTable, verifyKeyHashes, and FormulaKernelAdapter --- ...crementalChunkedCrossJoinStateManager.java | 16 ++++-- ...crementalChunkedCrossJoinStateManager.java | 16 ++++-- .../StaticChunkedCrossJoinStateManager.java | 16 ++++-- .../table/impl/SymbolTableCombiner.java | 16 ++++-- .../engine/table/impl/TimeTable.java | 16 +++--- .../select/formula/FormulaKernelAdapter.java | 55 +++++++++++-------- 6 files changed, 85 insertions(+), 50 deletions(-) diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/LeftOnlyIncrementalChunkedCrossJoinStateManager.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/LeftOnlyIncrementalChunkedCrossJoinStateManager.java index 2ffd6aca2bb..68f82fc56cf 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/LeftOnlyIncrementalChunkedCrossJoinStateManager.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/LeftOnlyIncrementalChunkedCrossJoinStateManager.java @@ -1175,20 +1175,26 @@ public boolean rehashRequired() { private void verifyKeyHashes() { final int maxSize = tableHashPivot; - final ChunkSource.FillContext [] keyFillContext = makeFillContexts(keySources, SharedContext.makeSharedContext(), maxSize); + final ChunkSource.FillContext [] keyFillContext = new ChunkSource.FillContext[keySources.length]; final WritableChunk [] keyChunks = getWritableKeyChunks(maxSize); - try (final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); + try (final SharedContext sharedContext = SharedContext.makeSharedContext(); + final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); + final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); final WritableBooleanChunk exists = WritableBooleanChunk.makeWritableChunk(maxSize); final WritableIntChunk hashChunk = WritableIntChunk.makeWritableChunk(maxSize); final WritableLongChunk tableLocationsChunk = WritableLongChunk.makeWritableChunk(maxSize); - final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); - final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final RowSet flatRowSet = RowSetFactory.flat(tableHashPivot); // @StateChunkName@ from \QObjectChunk\E final WritableObjectChunk stateChunk = WritableObjectChunk.makeWritableChunk(maxSize); final ChunkSource.FillContext fillContext = rightRowSetSource.makeFillContext(maxSize)) { - rightRowSetSource.fillChunk(fillContext, stateChunk, RowSetFactory.flat(tableHashPivot)); + for (int ii = 0; ii < keySources.length; ++ii) { + keyFillContext[ii] = keySources[ii].makeFillContext(maxSize, sharedContext); + } + + rightRowSetSource.fillChunk(fillContext, stateChunk, flatRowSet); ChunkUtils.fillInOrder(positions); diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/RightIncrementalChunkedCrossJoinStateManager.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/RightIncrementalChunkedCrossJoinStateManager.java index 955fdacef66..d5d1f7adf12 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/RightIncrementalChunkedCrossJoinStateManager.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/RightIncrementalChunkedCrossJoinStateManager.java @@ -1532,20 +1532,26 @@ public boolean rehashRequired() { private void verifyKeyHashes() { final int maxSize = tableHashPivot; - final ChunkSource.FillContext [] keyFillContext = makeFillContexts(keySources, SharedContext.makeSharedContext(), maxSize); + final ChunkSource.FillContext [] keyFillContext = new ChunkSource.FillContext[keySources.length]; final WritableChunk [] keyChunks = getWritableKeyChunks(maxSize); - try (final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); + try (final SharedContext sharedContext = SharedContext.makeSharedContext(); + final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); + final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); final WritableBooleanChunk exists = WritableBooleanChunk.makeWritableChunk(maxSize); final WritableIntChunk hashChunk = WritableIntChunk.makeWritableChunk(maxSize); final WritableLongChunk tableLocationsChunk = WritableLongChunk.makeWritableChunk(maxSize); - final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); - final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final RowSet flatRowSet = RowSetFactory.flat(tableHashPivot); // @StateChunkName@ from \QObjectChunk\E final WritableObjectChunk stateChunk = WritableObjectChunk.makeWritableChunk(maxSize); final ChunkSource.FillContext fillContext = rightRowSetSource.makeFillContext(maxSize)) { - rightRowSetSource.fillChunk(fillContext, stateChunk, RowSetFactory.flat(tableHashPivot)); + for (int ii = 0; ii < keySources.length; ++ii) { + keyFillContext[ii] = keySources[ii].makeFillContext(maxSize, sharedContext); + } + + rightRowSetSource.fillChunk(fillContext, stateChunk, flatRowSet); ChunkUtils.fillInOrder(positions); diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/StaticChunkedCrossJoinStateManager.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/StaticChunkedCrossJoinStateManager.java index 0d6771ea06d..5f26b6f2cb3 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/StaticChunkedCrossJoinStateManager.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/StaticChunkedCrossJoinStateManager.java @@ -1136,20 +1136,26 @@ public boolean rehashRequired() { private void verifyKeyHashes() { final int maxSize = tableHashPivot; - final ChunkSource.FillContext [] keyFillContext = makeFillContexts(keySources, SharedContext.makeSharedContext(), maxSize); + final ChunkSource.FillContext [] keyFillContext = new ChunkSource.FillContext[keySources.length]; final WritableChunk [] keyChunks = getWritableKeyChunks(maxSize); - try (final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); + try (final SharedContext sharedContext = SharedContext.makeSharedContext(); + final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); + final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); final WritableBooleanChunk exists = WritableBooleanChunk.makeWritableChunk(maxSize); final WritableIntChunk hashChunk = WritableIntChunk.makeWritableChunk(maxSize); final WritableLongChunk tableLocationsChunk = WritableLongChunk.makeWritableChunk(maxSize); - final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); - final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final RowSet flatRowSet = RowSetFactory.flat(tableHashPivot); // @StateChunkName@ from \QLongChunk\E final WritableLongChunk stateChunk = WritableLongChunk.makeWritableChunk(maxSize); final ChunkSource.FillContext fillContext = slotSource.makeFillContext(maxSize)) { - slotSource.fillChunk(fillContext, stateChunk, RowSetFactory.flat(tableHashPivot)); + for (int ii = 0; ii < keySources.length; ++ii) { + keyFillContext[ii] = keySources[ii].makeFillContext(maxSize, sharedContext); + } + + slotSource.fillChunk(fillContext, stateChunk, flatRowSet); ChunkUtils.fillInOrder(positions); diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/SymbolTableCombiner.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/SymbolTableCombiner.java index c4fbed83ffb..5312b7c13c1 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/SymbolTableCombiner.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/SymbolTableCombiner.java @@ -1006,20 +1006,26 @@ public boolean rehashRequired() { private void verifyKeyHashes() { final int maxSize = tableHashPivot; - final ChunkSource.FillContext [] keyFillContext = makeFillContexts(keySources, SharedContext.makeSharedContext(), maxSize); + final ChunkSource.FillContext [] keyFillContext = new ChunkSource.FillContext[keySources.length]; final WritableChunk [] keyChunks = getWritableKeyChunks(maxSize); - try (final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); + try (final SharedContext sharedContext = SharedContext.makeSharedContext(); + final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); + final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final WritableLongChunk positions = WritableLongChunk.makeWritableChunk(maxSize); final WritableBooleanChunk exists = WritableBooleanChunk.makeWritableChunk(maxSize); final WritableIntChunk hashChunk = WritableIntChunk.makeWritableChunk(maxSize); final WritableLongChunk tableLocationsChunk = WritableLongChunk.makeWritableChunk(maxSize); - final SafeCloseableArray ignored = new SafeCloseableArray<>(keyFillContext); - final SafeCloseableArray ignored2 = new SafeCloseableArray<>(keyChunks); + final RowSet flatRowSet = RowSetFactory.flat(tableHashPivot); // @StateChunkName@ from \QIntChunk\E final WritableIntChunk stateChunk = WritableIntChunk.makeWritableChunk(maxSize); final ChunkSource.FillContext fillContext = uniqueIdentifierSource.makeFillContext(maxSize)) { - uniqueIdentifierSource.fillChunk(fillContext, stateChunk, RowSetFactory.flat(tableHashPivot)); + for (int ii = 0; ii < keySources.length; ++ii) { + keyFillContext[ii] = keySources[ii].makeFillContext(maxSize, sharedContext); + } + + uniqueIdentifierSource.fillChunk(fillContext, stateChunk, flatRowSet); ChunkUtils.fillInOrder(positions); diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/TimeTable.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/TimeTable.java index c7abfce8798..0e34b11ccf6 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/TimeTable.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/TimeTable.java @@ -205,12 +205,7 @@ private void refresh(final boolean notifyListeners) { final boolean rowsAdded = rangeStart <= lastIndex; final boolean rowsRemoved = isBlinkTable && getRowSet().isNonempty(); if (rowsAdded || rowsRemoved) { - final RowSet addedRange = rowsAdded - ? RowSetFactory.fromRange(rangeStart, lastIndex) - : RowSetFactory.empty(); - final RowSet removedRange = rowsRemoved - ? RowSetFactory.fromRange(getRowSet().firstRowKey(), rangeStart - 1) - : RowSetFactory.empty(); + final long firstRemovedRowKey = rowsRemoved ? getRowSet().firstRowKey() : RowSequence.NULL_ROW_KEY; if (rowsAdded) { getRowSet().writableCast().insertRange(rangeStart, lastIndex); } @@ -218,7 +213,14 @@ private void refresh(final boolean notifyListeners) { getRowSet().writableCast().removeRange(0, rangeStart - 1); } if (notifyListeners) { - notifyListeners(addedRange, removedRange, RowSetFactory.empty()); + notifyListeners( + rowsAdded + ? RowSetFactory.fromRange(rangeStart, lastIndex) + : RowSetFactory.empty(), + rowsRemoved + ? RowSetFactory.fromRange(firstRemovedRowKey, rangeStart - 1) + : RowSetFactory.empty(), + RowSetFactory.empty()); } } } diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/select/formula/FormulaKernelAdapter.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/select/formula/FormulaKernelAdapter.java index 44c49fc521a..22753fd244c 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/select/formula/FormulaKernelAdapter.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/select/formula/FormulaKernelAdapter.java @@ -12,6 +12,7 @@ import io.deephaven.engine.rowset.RowSetFactory; import io.deephaven.engine.rowset.RowSequence; import io.deephaven.engine.rowset.TrackingRowSet; +import io.deephaven.util.SafeCloseableList; import io.deephaven.util.type.TypeUtils; import org.jetbrains.annotations.NotNull; @@ -306,33 +307,41 @@ public AdapterContext makeFillContext(final int chunkCapacity) { WritableLongChunk kChunk = null; final String[] sources = sourceDescriptor.sources; - // Create whichever of the three special chunks we need - for (final String source : sources) { - switch (source) { - case "i": - iChunk = WritableIntChunk.makeWritableChunk(chunkCapacity); - break; - case "ii": - iiChunk = WritableLongChunk.makeWritableChunk(chunkCapacity); - break; - case "k": - kChunk = WritableLongChunk.makeWritableChunk(chunkCapacity); - break; + // Everything allocated here is owned by partiallyBuilt until the AdapterContext exists to take ownership, so + // that a failure part way through does not orphan the chunks and contexts already created. + try (final SafeCloseableList partiallyBuilt = new SafeCloseableList()) { + // Create whichever of the three special chunks we need + for (final String source : sources) { + switch (source) { + case "i": + iChunk = partiallyBuilt.add(WritableIntChunk.makeWritableChunk(chunkCapacity)); + break; + case "ii": + iiChunk = partiallyBuilt.add(WritableLongChunk.makeWritableChunk(chunkCapacity)); + break; + case "k": + kChunk = partiallyBuilt.add(WritableLongChunk.makeWritableChunk(chunkCapacity)); + break; + } } - } - // Make contexts -- we leave nulls in the slots where i, ii, or k would be. - final ColumnSource.GetContext[] sourceContexts = new ColumnSource.GetContext[sources.length]; - for (int ii = 0; ii < sources.length; ++ii) { - final String name = sources[ii]; - if (name.equals("i") || name.equals("ii") || name.equals("k")) { - continue; + // Make contexts -- we leave nulls in the slots where i, ii, or k would be. + final ColumnSource.GetContext[] sourceContexts = + partiallyBuilt.addArray(new ColumnSource.GetContext[sources.length]); + for (int ii = 0; ii < sources.length; ++ii) { + final String name = sources[ii]; + if (name.equals("i") || name.equals("ii") || name.equals("k")) { + continue; + } + final ColumnSource cs = columnSources.get(name); + sourceContexts[ii] = cs.makeGetContext(chunkCapacity); } - final ColumnSource cs = columnSources.get(name); - sourceContexts[ii] = cs.makeGetContext(chunkCapacity); + final FillContext kernelContext = partiallyBuilt.add(kernel.makeFillContext(chunkCapacity)); + final AdapterContext result = + new AdapterContext(iChunk, iiChunk, kChunk, sourceContexts, kernelContext); + partiallyBuilt.clear(); + return result; } - final FillContext kernelContext = kernel.makeFillContext(chunkCapacity); - return new AdapterContext(iChunk, iiChunk, kChunk, sourceContexts, kernelContext); } private static class AdapterContext implements FillContext { From 8429806d805e037f31053e19eafced8939c8335a Mon Sep 17 00:00:00 2001 From: "Charles P. Wright" Date: Thu, 27 Aug 2026 08:28:54 -0400 Subject: [PATCH 2/3] fix: DH-23494: Close per-cycle row set leaks in slice, cross join, natural join, as-of join, snapshot, and filters --- .../engine/table/impl/AsOfJoinHelper.java | 14 +-- .../impl/CrossJoinModifiedSlotTracker.java | 91 ++++++++++--------- .../engine/table/impl/SliceLikeOperation.java | 46 +++++----- .../table/impl/TableUpdateValidator.java | 50 ++++++---- ...entalNaturalJoinStateManagerTypedBase.java | 6 ++ ...entalNaturalJoinStateManagerTypedBase.java | 6 ++ .../snapshot/SnapshotIncrementalListener.java | 7 +- .../table/impl/util/SyncTableFilter.java | 27 ++++-- .../engine/util/LeaderTableFilter.java | 13 ++- 9 files changed, 155 insertions(+), 105 deletions(-) diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/AsOfJoinHelper.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/AsOfJoinHelper.java index c6d1cdcc23a..59484f3d31f 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/AsOfJoinHelper.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/AsOfJoinHelper.java @@ -760,14 +760,16 @@ public void onUpdate(TableUpdate upstream) { final SegmentedSortedArray rightSsa = asOfJoinStateManager.getRightSsa(slot); - final RowSet rightAdded = sequentialBuilders.get(slotIndex).build(); - sequentialBuilders.set(slotIndex, null); + final int rightSize; + try (final RowSet rightAdded = sequentialBuilders.get(slotIndex).build()) { + sequentialBuilders.set(slotIndex, null); - final int rightSize = rightAdded.intSize(); + rightSize = rightAdded.intSize(); - rightStampSource.fillChunk(rightStampFillContext.ensureCapacity(rightSize), - rightStampChunk.ensureCapacity(rightSize), rightAdded); - rightAdded.fillRowKeyChunk(insertedIndices.ensureCapacity(rightSize)); + rightStampSource.fillChunk(rightStampFillContext.ensureCapacity(rightSize), + rightStampChunk.ensureCapacity(rightSize), rightAdded); + rightAdded.fillRowKeyChunk(insertedIndices.ensureCapacity(rightSize)); + } sortContext.ensureCapacity(rightSize).sort(insertedIndices.get(), rightStampChunk.get()); final int valuesWithNext = rightSsa.insertAndGetNextValue(rightStampChunk.get(), diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/CrossJoinModifiedSlotTracker.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/CrossJoinModifiedSlotTracker.java index f62eb8da90e..f6fd0da661b 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/CrossJoinModifiedSlotTracker.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/CrossJoinModifiedSlotTracker.java @@ -507,20 +507,21 @@ void flushLeftRemoves() { continue; } - final RowSet leftRemoved = slotState.indexBuilder.build(); - slotState.leftRowSet.remove(leftRemoved); - slotState.indexBuilder = RowSetFactory.builderRandom(); - long sizePrev = slotState.rightRowSet.sizePrev(); - if (sizePrev > 0) { - leftRemoved.forAllRowKeys(ii -> { - final long prevOffset = ii << jsm.getPrevNumShiftBits(); - builder.addRange(prevOffset, prevOffset + sizePrev - 1); - }); - } else if (jsm.leftOuterJoin()) { - leftRemoved.forAllRowKeys(ii -> { - final long prevOffset = ii << jsm.getPrevNumShiftBits(); - builder.addKey(prevOffset); - }); + final long sizePrev = slotState.rightRowSet.sizePrev(); + try (final RowSet slotLeftRemoved = slotState.indexBuilder.build()) { + slotState.leftRowSet.remove(slotLeftRemoved); + slotState.indexBuilder = RowSetFactory.builderRandom(); + if (sizePrev > 0) { + slotLeftRemoved.forAllRowKeys(ii -> { + final long prevOffset = ii << jsm.getPrevNumShiftBits(); + builder.addRange(prevOffset, prevOffset + sizePrev - 1); + }); + } else if (jsm.leftOuterJoin()) { + slotLeftRemoved.forAllRowKeys(ii -> { + final long prevOffset = ii << jsm.getPrevNumShiftBits(); + builder.addKey(prevOffset); + }); + } } } leftRemoved = builder.build(); @@ -534,22 +535,23 @@ void flushLeftAdds() { continue; } - final RowSet leftAdded = slotState.indexBuilder.build(); - slotState.leftRowSet.insert(leftAdded); - jsm.updateLeftRowRedirection(leftAdded, slotState.slotLocation); - - slotState.indexBuilder = null; - long size = slotState.rightRowSet.size(); - if (size > 0) { - leftAdded.forAllRowKeys(ii -> { - final long currOffset = ii << jsm.getNumShiftBits(); - downstreamAdds.addRange(currOffset, currOffset + size - 1); - }); - } else if (jsm.leftOuterJoin()) { - leftAdded.forAllRowKeys(ii -> { - final long currOffset = ii << jsm.getNumShiftBits(); - downstreamAdds.addKey(currOffset); - }); + final long size = slotState.rightRowSet.size(); + try (final RowSet slotLeftAdded = slotState.indexBuilder.build()) { + slotState.leftRowSet.insert(slotLeftAdded); + jsm.updateLeftRowRedirection(slotLeftAdded, slotState.slotLocation); + + slotState.indexBuilder = null; + if (size > 0) { + slotLeftAdded.forAllRowKeys(ii -> { + final long currOffset = ii << jsm.getNumShiftBits(); + downstreamAdds.addRange(currOffset, currOffset + size - 1); + }); + } else if (jsm.leftOuterJoin()) { + slotLeftAdded.forAllRowKeys(ii -> { + final long currOffset = ii << jsm.getNumShiftBits(); + downstreamAdds.addKey(currOffset); + }); + } } final RowSetBuilderRandom modifiedAdds = RowSetFactory.builderRandom(); @@ -598,21 +600,22 @@ void flushLeftModifies() { if (slotState == null) { continue; } - final RowSet leftRemoved = slotState.indexBuilder.build(); - slotState.leftRowSet.remove(leftRemoved); - jsm.updateLeftRowRedirection(leftRemoved, RowSequence.NULL_ROW_KEY); - slotState.indexBuilder = RowSetFactory.builderRandom(); final long sizePrev = slotState.rightRowSet.sizePrev(); - if (sizePrev > 0) { - leftRemoved.forAllRowKeys(ii -> { - final long prevOffset = ii << jsm.getPrevNumShiftBits(); - rmBuilder.addRange(prevOffset, prevOffset + sizePrev - 1); - }); - } else if (jsm.leftOuterJoin()) { - leftRemoved.forAllRowKeys(ii -> { - final long prevOffset = ii << jsm.getPrevNumShiftBits(); - rmBuilder.addKey(prevOffset); - }); + try (final RowSet slotLeftRemoved = slotState.indexBuilder.build()) { + slotState.leftRowSet.remove(slotLeftRemoved); + jsm.updateLeftRowRedirection(slotLeftRemoved, RowSequence.NULL_ROW_KEY); + slotState.indexBuilder = RowSetFactory.builderRandom(); + if (sizePrev > 0) { + slotLeftRemoved.forAllRowKeys(ii -> { + final long prevOffset = ii << jsm.getPrevNumShiftBits(); + rmBuilder.addRange(prevOffset, prevOffset + sizePrev - 1); + }); + } else if (jsm.leftOuterJoin()) { + slotLeftRemoved.forAllRowKeys(ii -> { + final long prevOffset = ii << jsm.getPrevNumShiftBits(); + rmBuilder.addKey(prevOffset); + }); + } } final long size = slotState.rightRowSet.size(); if (sizePrev > 0 && size > 0) { diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/SliceLikeOperation.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/SliceLikeOperation.java index 1b96a647219..a7f04a68e90 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/SliceLikeOperation.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/SliceLikeOperation.java @@ -157,35 +157,37 @@ public void onUpdate(TableUpdate upstream) { private void onUpdate(final TableUpdate upstream) { final TrackingWritableRowSet rowSet = resultTable.getRowSet().writableCast(); - final RowSet sliceRowSet = computeSliceRowSet(parent.getRowSet()); - final TableUpdateImpl downstream = new TableUpdateImpl(); - downstream.removed = upstream.removed().intersect(rowSet); - rowSet.remove(downstream.removed()); + try (final RowSet sliceRowSet = computeSliceRowSet(parent.getRowSet())) { + final TableUpdateImpl downstream = new TableUpdateImpl(); + downstream.removed = upstream.removed().intersect(rowSet); + rowSet.remove(downstream.removed()); - downstream.shifted = upstream.shifted().intersect(rowSet); - downstream.shifted().apply(rowSet); + downstream.shifted = upstream.shifted().intersect(rowSet); + downstream.shifted().apply(rowSet); - // Must calculate in post-shift space what indices were removed by the slice operation. - final WritableRowSet opRemoved = rowSet.minus(sliceRowSet); - rowSet.remove(opRemoved); - downstream.shifted().unapply(opRemoved); - downstream.removed().writableCast().insert(opRemoved); + // Must calculate in post-shift space what indices were removed by the slice operation. + try (final WritableRowSet opRemoved = rowSet.minus(sliceRowSet)) { + rowSet.remove(opRemoved); + downstream.shifted().unapply(opRemoved); + downstream.removed().writableCast().insert(opRemoved); + } - // Must intersect against modified set before adding the new rows to result rowSet. - downstream.modified = upstream.modified().intersect(rowSet); + // Must intersect against modified set before adding the new rows to result rowSet. + downstream.modified = upstream.modified().intersect(rowSet); - downstream.added = sliceRowSet.minus(rowSet); - rowSet.insert(downstream.added()); + downstream.added = sliceRowSet.minus(rowSet); + rowSet.insert(downstream.added()); - // propagate an empty MCS if modified is empty - downstream.modifiedColumnSet = upstream.modifiedColumnSet(); - if (downstream.modified().isEmpty()) { - downstream.modifiedColumnSet = resultTable.getModifiedColumnSetForUpdates(); - downstream.modifiedColumnSet.clear(); - } + // propagate an empty MCS if modified is empty + downstream.modifiedColumnSet = upstream.modifiedColumnSet(); + if (downstream.modified().isEmpty()) { + downstream.modifiedColumnSet = resultTable.getModifiedColumnSetForUpdates(); + downstream.modifiedColumnSet.clear(); + } - resultTable.notifyListeners(downstream); + resultTable.notifyListeners(downstream); + } } private WritableRowSet computeSliceRowSet(RowSet useRowSet) { diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/TableUpdateValidator.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/TableUpdateValidator.java index 9de98526f9c..8caa2630cee 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/TableUpdateValidator.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/TableUpdateValidator.java @@ -180,7 +180,9 @@ private void onUpdate(final TableUpdate upstream) { false); } - validateIndexesEqual("pre-update rowSet", rowSet, tableToValidate.getRowSet().copyPrev()); + try (final RowSet prevRowSet = tableToValidate.getRowSet().copyPrev()) { + validateIndexesEqual("pre-update rowSet", rowSet, prevRowSet); + } rowSet.remove(upstream.removed()); Arrays.stream(columnInfos).forEach((ci) -> ci.remove(upstream.removed())); @@ -189,8 +191,9 @@ private void onUpdate(final TableUpdate upstream) { upstream.shifted().apply(rowSet); if (aggressiveUpdateValidation) { - final RowSet unmodified = rowSet.minus(upstream.modified()); - validateValues("post-shift unmodified", ModifiedColumnSet.ALL, unmodified, false, false); + try (final RowSet unmodified = rowSet.minus(upstream.modified())) { + validateValues("post-shift unmodified", ModifiedColumnSet.ALL, unmodified, false, false); + } validateValues("post-shift unmodified columns", upstream.modifiedColumnSet(), upstream.modified(), false, true); @@ -198,8 +201,11 @@ private void onUpdate(final TableUpdate upstream) { // added if (rowSet.overlaps(upstream.added())) { - noteIssue(() -> "post-shift rowSet contains rows that are added: " - + rowSet.intersect(upstream.added())); + noteIssue(() -> { + try (final RowSet addedInRowSet = rowSet.intersect(upstream.added())) { + return "post-shift rowSet contains rows that are added: " + addedInRowSet; + } + }); } rowSet.insert(upstream.added()); validateIndexesEqual("post-update rowSet", rowSet, tableToValidate.getRowSet()); @@ -208,12 +214,19 @@ private void onUpdate(final TableUpdate upstream) { // modified updateValues(upstream.modifiedColumnSet(), upstream.modified(), false); if (upstream.added().overlaps(upstream.modified())) { - noteIssue(() -> "added contains rows that are modified (post-shift): " - + upstream.added().intersect(upstream.modified())); + noteIssue(() -> { + try (final RowSet addedAndModified = upstream.added().intersect(upstream.modified())) { + return "added contains rows that are modified (post-shift): " + addedAndModified; + } + }); } if (upstream.removed().overlaps(upstream.getModifiedPreShift())) { - noteIssue(() -> "removed contains rows that are modified (pre-shift): " - + upstream.removed().intersect(upstream.getModifiedPreShift())); + noteIssue(() -> { + try (final RowSet removedAndModified = + upstream.removed().intersect(upstream.getModifiedPreShift())) { + return "removed contains rows that are modified (pre-shift): " + removedAndModified; + } + }); } if (!issues.isEmpty()) { @@ -240,13 +253,14 @@ private void validateIndexesEqual(final String what, final RowSet expected, fina return; } - final RowSet missing = expected.minus(actual); - final RowSet excess = actual.minus(expected); - if (missing.isNonempty()) { - noteIssue(() -> what + " expected.minus(actual)=" + missing); - } - if (excess.isNonempty()) { - noteIssue(() -> what + " actual.minus(expected)=" + excess); + try (final RowSet missing = expected.minus(actual); + final RowSet excess = actual.minus(expected)) { + if (missing.isNonempty()) { + noteIssue(() -> what + " expected.minus(actual)=" + missing); + } + if (excess.isNonempty()) { + noteIssue(() -> what + " actual.minus(expected)=" + excess); + } } } @@ -417,7 +431,9 @@ private WritableBooleanChunk equalValuesDest() { @Override public void shift(final long beginRange, final long endRange, final long shiftDelta) { - ((RowSetShiftCallback) expectedSource).shift(rowSet.subSetByKeyRange(beginRange, endRange), shiftDelta); + try (final RowSet keysToShift = rowSet.subSetByKeyRange(beginRange, endRange)) { + ((RowSetShiftCallback) expectedSource).shift(keysToShift, shiftDelta); + } } public void remove(final RowSet toRemove) { diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/IncrementalNaturalJoinStateManagerTypedBase.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/IncrementalNaturalJoinStateManagerTypedBase.java index 000d663b9e0..3ded542b730 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/IncrementalNaturalJoinStateManagerTypedBase.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/IncrementalNaturalJoinStateManagerTypedBase.java @@ -440,6 +440,12 @@ protected long allocateDuplicateLocation() { } protected void freeDuplicateLocation(long duplicateLocation) { + // The duplicate row set at this location is no longer reachable; the location is reused by + // allocateDuplicateLocation, which overwrites the slot with a freshly built row set. + final WritableRowSet duplicates = rightSideDuplicateRowSets.getAndSetUnsafe(duplicateLocation, null); + if (duplicates != null) { + duplicates.close(); + } freeDuplicateValues.add(duplicateLocation); } diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/RightIncrementalNaturalJoinStateManagerTypedBase.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/RightIncrementalNaturalJoinStateManagerTypedBase.java index 20eaf567fee..0e7536bbe33 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/RightIncrementalNaturalJoinStateManagerTypedBase.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/naturaljoin/RightIncrementalNaturalJoinStateManagerTypedBase.java @@ -217,6 +217,12 @@ protected long allocateDuplicateLocation() { } protected void freeDuplicateLocation(long duplicateLocation) { + // The duplicate row set at this location is no longer reachable; the location is reused by + // allocateDuplicateLocation, which overwrites the slot with a freshly built row set. + final WritableRowSet duplicates = rightSideDuplicateRowSets.getAndSetUnsafe(duplicateLocation, null); + if (duplicates != null) { + duplicates.close(); + } freeDuplicateValues.add(duplicateLocation); } diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/snapshot/SnapshotIncrementalListener.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/snapshot/SnapshotIncrementalListener.java index 2906c460b74..e3bdef60e0a 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/snapshot/SnapshotIncrementalListener.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/snapshot/SnapshotIncrementalListener.java @@ -92,9 +92,12 @@ public void doSnapshot() { final RowSet baseAdded = expander.getAdded().copy(); final RowSet baseModified = expander.getModified().copy(); final RowSet baseRemoved = expander.getRemoved().copy(); - final RowSet rowsToCopy = baseAdded.union(baseModified); - doRowCopy(rowsToCopy); + // baseAdded, baseModified, and baseRemoved are given away to notifyListeners below; rowsToCopy exists + // only to drive the copy. + try (final RowSet rowsToCopy = baseAdded.union(baseModified)) { + doRowCopy(rowsToCopy); + } resultTable.getRowSet().writableCast().update(baseAdded, baseRemoved); resultTable.notifyListeners(baseAdded, baseRemoved, baseModified); diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/util/SyncTableFilter.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/util/SyncTableFilter.java index ac073b48e79..0ec88c8ce09 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/util/SyncTableFilter.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/util/SyncTableFilter.java @@ -24,6 +24,7 @@ import io.deephaven.engine.table.impl.TupleSourceFactory; import io.deephaven.engine.rowset.chunkattributes.OrderedRowKeys; import io.deephaven.util.QueryConstants; +import io.deephaven.util.SafeCloseable; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; @@ -198,8 +199,9 @@ protected void process() { if (recorder.getShifted().nonempty()) { throw new IllegalStateException("Can not process shifted rows in SyncTableFilter!"); } - final RowSet addedAndModified = recorder.getAdded().union(recorder.getModified()); - consumeRows(rr, addedAndModified); + try (final RowSet addedAndModified = recorder.getAdded().union(recorder.getModified())) { + consumeRows(rr, addedAndModified); + } } } @@ -224,10 +226,11 @@ protected void process() { if (!keysToRefilter.contains(key)) { // if we did not refilter this key; then we should add the currently matched values, // otherwise we ignore them because they have already been superseded - final WritableRowSet newlyMatchedRows = state.currentIdBuilder.build(); - state.matchedRows.insert(newlyMatchedRows); - newlyMatchedRows.remove(resultRowSet[tt]); - addedBuilder.addRowSet(newlyMatchedRows); + try (final WritableRowSet newlyMatchedRows = state.currentIdBuilder.build()) { + state.matchedRows.insert(newlyMatchedRows); + newlyMatchedRows.remove(resultRowSet[tt]); + addedBuilder.addRowSet(newlyMatchedRows); + } } state.currentIdBuilder = null; } @@ -237,19 +240,23 @@ protected void process() { resultRowSet[tt].remove(removed); resultRowSet[tt].insert(added); - final WritableRowSet addedAndRemoved = added.intersect(removed); - final WritableRowSet modified; if (recorders.get(tt).getNotificationStep() == currentStep) { modified = recorders.get(tt).getModified().intersect(resultRowSet[tt]); modified.remove(added); - modified.insert(addedAndRemoved); + try (final WritableRowSet addedAndRemoved = added.intersect(removed)) { + modified.insert(addedAndRemoved); + } } else { - modified = addedAndRemoved; + // Rows that were both added and removed are the only modifications in this case. + modified = added.intersect(removed); } if (added.isNonempty() || removed.isNonempty() || modified.isNonempty()) { results[tt].notifyListeners(added, removed, modified); + } else { + // There is nothing to notify, so nothing takes ownership of these. + SafeCloseable.closeAll(added, removed, modified); } } keysToRefilter.clear(); diff --git a/engine/table/src/main/java/io/deephaven/engine/util/LeaderTableFilter.java b/engine/table/src/main/java/io/deephaven/engine/util/LeaderTableFilter.java index e778f085ba4..2830a4f42f9 100644 --- a/engine/table/src/main/java/io/deephaven/engine/util/LeaderTableFilter.java +++ b/engine/table/src/main/java/io/deephaven/engine/util/LeaderTableFilter.java @@ -30,6 +30,7 @@ import io.deephaven.engine.updategraph.UpdateGraph; import io.deephaven.engine.util.systemicmarking.SystemicObjectTracker; import io.deephaven.util.QueryConstants; +import io.deephaven.util.SafeCloseable; import io.deephaven.util.SafeCloseableArray; import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.Nullable; @@ -483,10 +484,11 @@ protected void process() { if (!processPendingResult.keysToRefilter.contains(key)) { // if we did not refilter this key; then we should add the currently matched values, // otherwise we ignore them because they have already been superseded - final WritableRowSet newlyMatchedRows = state.currentIdBuilder.build(); - state.matchedRows.insert(newlyMatchedRows); - newlyMatchedRows.remove(followerResultRowSets[tt]); - addedBuilder.addRowSet(newlyMatchedRows); + try (final WritableRowSet newlyMatchedRows = state.currentIdBuilder.build()) { + state.matchedRows.insert(newlyMatchedRows); + newlyMatchedRows.remove(followerResultRowSets[tt]); + addedBuilder.addRowSet(newlyMatchedRows); + } } state.currentIdBuilder = null; } @@ -507,6 +509,9 @@ protected void process() { update.modifiedColumnSet = ModifiedColumnSet.EMPTY; update.shifted = RowSetShiftData.EMPTY; followerResults[tt].notifyListeners(update); + } else { + // There is nothing to notify, so no update takes ownership of these. + SafeCloseable.closeAll(added, removed); } } From be79f1af2f45ca0f0c995ab46ce6309f521cfc7f Mon Sep 17 00:00:00 2001 From: "Charles P. Wright" Date: Thu, 27 Aug 2026 09:08:20 -0400 Subject: [PATCH 3/3] fix: DH-23494: Close SparseArraySource shift iterators --- .../impl/sources/BooleanSparseArraySource.java | 17 ++++++++++------- .../impl/sources/ByteSparseArraySource.java | 17 ++++++++++------- .../sources/CharacterSparseArraySource.java | 17 ++++++++++------- .../impl/sources/DoubleSparseArraySource.java | 17 ++++++++++------- .../impl/sources/FloatSparseArraySource.java | 17 ++++++++++------- .../impl/sources/IntegerSparseArraySource.java | 17 ++++++++++------- .../impl/sources/LongSparseArraySource.java | 17 ++++++++++------- .../impl/sources/ObjectSparseArraySource.java | 17 ++++++++++------- .../impl/sources/ShortSparseArraySource.java | 17 ++++++++++------- 9 files changed, 90 insertions(+), 63 deletions(-) diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/BooleanSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/BooleanSparseArraySource.java index 6b06ba8d4cc..4da9c9eeaed 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/BooleanSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/BooleanSparseArraySource.java @@ -141,13 +141,16 @@ public final void set(long key, byte value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getBoolean(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getBoolean(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ByteSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ByteSparseArraySource.java index a80f478a5dc..675f5490053 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ByteSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ByteSparseArraySource.java @@ -135,13 +135,16 @@ public final void set(long key, byte value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getByte(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getByte(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/CharacterSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/CharacterSparseArraySource.java index 42ac6442f5c..78be16c009d 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/CharacterSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/CharacterSparseArraySource.java @@ -131,13 +131,16 @@ public final void set(long key, char value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getChar(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getChar(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/DoubleSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/DoubleSparseArraySource.java index c4016316035..fd96b5f5386 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/DoubleSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/DoubleSparseArraySource.java @@ -135,13 +135,16 @@ public final void set(long key, double value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getDouble(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getDouble(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/FloatSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/FloatSparseArraySource.java index 15db86076ac..596479b6849 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/FloatSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/FloatSparseArraySource.java @@ -135,13 +135,16 @@ public final void set(long key, float value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getFloat(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getFloat(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/IntegerSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/IntegerSparseArraySource.java index 7f350e289c3..40d3638ae21 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/IntegerSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/IntegerSparseArraySource.java @@ -135,13 +135,16 @@ public final void set(long key, int value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getInt(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getInt(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/LongSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/LongSparseArraySource.java index 35d48563096..2cfcb4b0b76 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/LongSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/LongSparseArraySource.java @@ -143,13 +143,16 @@ public final void set(long key, long value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getLong(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getLong(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ObjectSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ObjectSparseArraySource.java index 82bf2301fca..4e05d32eb17 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ObjectSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ObjectSparseArraySource.java @@ -134,13 +134,16 @@ public final void set(long key, T value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, get(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, get(i)); + setNull(i); + return true; + }); + } } // region boxed methods diff --git a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ShortSparseArraySource.java b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ShortSparseArraySource.java index ff206d3cf33..aa1ee567b39 100644 --- a/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ShortSparseArraySource.java +++ b/engine/table/src/main/java/io/deephaven/engine/table/impl/sources/ShortSparseArraySource.java @@ -135,13 +135,16 @@ public final void set(long key, short value) { @Override public void shift(final RowSet keysToShift, final long shiftDelta) { - final RowSet.SearchIterator it = - (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator(); - it.forEachLong((i) -> { - set(i + shiftDelta, getShort(i)); - setNull(i); - return true; - }); + // Shifting up has to walk the keys in descending order so that a shifted value never overwrites one that has + // not been read yet. + try (final RowSet.SearchIterator it = + (shiftDelta > 0) ? keysToShift.reverseIterator() : keysToShift.searchIterator()) { + it.forEachLong((i) -> { + set(i + shiftDelta, getShort(i)); + setNull(i); + return true; + }); + } } // region boxed methods