This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new d512b22ecbd Track segments with snapshot to improve upsert snapshot
taking flow (#19028)
d512b22ecbd is described below
commit d512b22ecbdf88d0fe4921b88435bad3c085b808
Author: Xiaotian (Jackie) Jiang <[email protected]>
AuthorDate: Wed Jul 22 13:52:03 2026 -0700
Track segments with snapshot to improve upsert snapshot taking flow (#19028)
---
.../upsert/BasePartitionUpsertMetadataManager.java | 206 +++++++++++++--------
.../BasePartitionUpsertMetadataManagerTest.java | 73 +++++++-
2 files changed, 199 insertions(+), 80 deletions(-)
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
index 733e4e24014..837ce9c6431 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
@@ -23,8 +23,8 @@ import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.AtomicDouble;
import java.io.File;
import java.io.IOException;
+import java.util.ArrayList;
import java.util.HashMap;
-import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
@@ -101,11 +101,18 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
// is not managed by the manager currently.
protected final Set<IndexSegment> _trackedSegments =
ConcurrentHashMap.newKeySet();
// Track all the immutable segments where changes took place since last
snapshot was taken.
- // Note: we need take to take _snapshotLock RLock while updating this set as
it may be updated by the multiple
- // Helix threads. Otherwise, segments might be missed by the consuming
thread when taking snapshots, which takes
- // snapshotLock WLock and clear the tracking set to avoid keeping segment
object references around.
+ // Note: the set can be updated by multiple Helix task threads and the
consuming thread concurrently. To not miss
+ // segments updated while a snapshot is being taken, the snapshot flow
removes the exact segments it has persisted
+ // instead of clearing the set, and it retains only the tracked segments at
the end of each snapshot round to avoid
+ // keeping stale segment object references around.
// Skip mutableSegments as only immutable segments are for taking snapshots.
- protected final Set<IndexSegment> _updatedSegmentsSinceLastSnapshot =
ConcurrentHashMap.newKeySet();
+ protected final Set<ImmutableSegmentImpl> _updatedSegmentsSinceLastSnapshot
= ConcurrentHashMap.newKeySet();
+ // Track all the immutable segments that have their validDocIds snapshot
file persisted on disk. The snapshot flow
+ // persists the snapshots of these segments before creating snapshot files
for the other segments, and uses this
+ // set to classify the segments without checking the segment directory or
acquiring the segmentLock. The
+ // queryableDocIds snapshot file is persisted along with the validDocIds
snapshot file, but doesn't decide the
+ // membership of this set. Only maintained when snapshot is enabled.
+ protected final Set<ImmutableSegmentImpl> _segmentsWithSnapshot =
ConcurrentHashMap.newKeySet();
// NOTE: We do not persist snapshot on the first consuming segment because
most segments might not be loaded yet
// We only do this for Full-Upsert tables, for partial-upsert tables, we
have a check allSegmentsLoaded
@@ -291,12 +298,11 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
_logger.info("Skip adding segment: {} because metadata manager is
already stopped", segment.getSegmentName());
return;
}
+ ImmutableSegmentImpl immutableSegment = (ImmutableSegmentImpl) segment;
try {
- doAddSegment((ImmutableSegmentImpl) segment);
- _trackedSegments.add(segment);
- if (_enableSnapshot) {
- _updatedSegmentsSinceLastSnapshot.add(segment);
- }
+ doAddSegment(immutableSegment);
+ _trackedSegments.add(immutableSegment);
+ trackSegmentForSnapshot(immutableSegment);
} finally {
finishOperation();
}
@@ -417,10 +423,11 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
_logger.info("Skip preloading segment: {} because metadata manager is
already stopped", segmentName);
return;
}
+ ImmutableSegmentImpl immutableSegment = (ImmutableSegmentImpl) segment;
try {
- doPreloadSegment((ImmutableSegmentImpl) segment);
- _trackedSegments.add(segment);
- _updatedSegmentsSinceLastSnapshot.add(segment);
+ doPreloadSegment(immutableSegment);
+ _trackedSegments.add(immutableSegment);
+ trackSegmentForSnapshot(immutableSegment);
} finally {
finishOperation();
}
@@ -552,8 +559,8 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
segment.getSegmentName());
return false;
}
- // NOTE: We don't acquire snapshot read lock here because snapshot is
always taken before a new consuming segment
- // starts consuming, so it won't overlap with this method
+ // NOTE: Snapshot is always taken by the consuming thread before a new
consuming segment starts consuming, so
+ // taking snapshot won't overlap with this method
try {
boolean addRecord = doAddRecord(segment, recordInfo);
_trackedSegments.add(segment);
@@ -577,15 +584,13 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
}
try {
doReplaceSegment(segment, oldSegment);
- if (!(segment instanceof EmptyIndexSegment)) {
- _trackedSegments.add(segment);
- if (_enableSnapshot) {
- _updatedSegmentsSinceLastSnapshot.add(segment);
- }
+ if (segment instanceof ImmutableSegmentImpl immutableSegment) {
+ _trackedSegments.add(immutableSegment);
+ // The snapshot file is evaluated on the new segment object because
the replacement might have kept it
+ // (e.g. segment reload) or started from a clean segment directory
(e.g. segment re-download).
+ trackSegmentForSnapshot(immutableSegment);
}
- _trackedSegments.remove(oldSegment);
- // Remove the segment eagerly so we need not wait until next snapshot
cycle to remove replaced segment
- _updatedSegmentsSinceLastSnapshot.remove(oldSegment);
+ untrackSegment(oldSegment);
} finally {
finishOperation();
}
@@ -617,9 +622,8 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
// build a fresh segment without one and fall through to the full scan
below, unaffected. Rebuilding from just the
// snapshot's valid docs avoids resurrecting primary keys already
expired/deleted by TTL.
MutableRoaringBitmap validDocIdsSnapshot = null;
- if (isTTLEnabled() && segment instanceof ImmutableSegmentImpl) {
- validDocIdsSnapshot =
- ((ImmutableSegmentImpl)
segment).loadDocIdsFromSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME);
+ if (isTTLEnabled() && segment instanceof ImmutableSegmentImpl
immutableSegment) {
+ validDocIdsSnapshot =
immutableSegment.loadDocIdsFromSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME);
}
Iterator<RecordInfo> recordInfoIterator;
if (validDocIdsSnapshot != null) {
@@ -713,7 +717,8 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
*/
public boolean shouldRevertMetadataOnInconsistency(IndexSegment oldSegment) {
return
ConsumingSegmentConsistencyModeListener.getInstance().getConsistencyMode()
- .equals(ConsumingSegmentConsistencyModeListener.Mode.PROTECTED) &&
oldSegment instanceof MutableSegment
+ == ConsumingSegmentConsistencyModeListener.Mode.PROTECTED
+ && oldSegment instanceof MutableSegment
&& _context.isTableTypeInconsistentDuringConsumption();
}
@@ -782,9 +787,7 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
} else {
doRemoveSegment(segment);
}
- _trackedSegments.remove(segment);
- // Remove the segment eagerly so we need not wait until next snapshot
cycle to remove dropped segment
- _updatedSegmentsSinceLastSnapshot.remove(segment);
+ untrackSegment(segment);
} finally {
finishOperation();
}
@@ -930,69 +933,76 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
int numImmutableSegments = 0;
int numConsumingSegments = 0;
int numUnchangedSegments = 0;
- // The segments without validDocIds & queryable docId snapshots should
take their snapshots at last. So that when
- // there is failure to take snapshots, the validDocIds snapshot on disk
still keep track of an exclusive set of
- // valid docs across segments. Because the valid docs as tracked by the
existing validDocIds snapshots can only
- // get less. That no overlap of valid docs among segments with snapshots
is required by the preloading to work
- // correctly. We wouldn't be using queryableDocIds anywhere currently
during preload - storing them so we
- // could better extend the functionality. Best case scenario, if both
validDocs and queryableDocs are not persisted,
- // we will be considering that segment is not storing updated bitmap
copies on disk
- Set<ImmutableSegmentImpl> segmentsWithoutSnapshot = new HashSet<>();
- TableDataManager tableDataManager = _context.getTableDataManager();
- Preconditions.checkNotNull(tableDataManager, "Taking snapshot requires
tableDataManager");
- boolean isSegmentSkipped = false;
+ // The segments without snapshots on disk should take their snapshots
after the segments with existing snapshots
+ // persist theirs, so that when there is failure to persist the existing
snapshots, no new snapshot file is added
+ // on disk, and the validDocIds snapshots kept on disk still track an
exclusive set of valid docs across
+ // segments. This works because the valid docs as tracked by the existing
snapshots can only get less, and the
+ // preloading requires no overlap of valid docs among the snapshots on
disk to work correctly. The
+ // queryableDocIds snapshots are not used during preload currently -
storing them so we could better extend the
+ // functionality.
+ List<ImmutableSegmentImpl> segmentsWithSnapshot = new ArrayList<>();
+ List<ImmutableSegmentImpl> segmentsWithoutSnapshot = new ArrayList<>();
for (IndexSegment segment : _trackedSegments) {
- if (!(segment instanceof ImmutableSegmentImpl)) {
+ if (segment instanceof ImmutableSegmentImpl immutableSegment) {
+ if (!_updatedSegmentsSinceLastSnapshot.contains(segment)) {
+ // if no updates since last snapshot then skip
+ numUnchangedSegments++;
+ continue;
+ }
+ if (_segmentsWithSnapshot.contains(segment)) {
+ segmentsWithSnapshot.add(immutableSegment);
+ } else {
+ segmentsWithoutSnapshot.add(immutableSegment);
+ }
+ } else {
numConsumingSegments++;
- continue;
- }
- if (!_updatedSegmentsSinceLastSnapshot.contains(segment)) {
- // if no updates since last snapshot then skip
- numUnchangedSegments++;
- continue;
}
- // Try to acquire the segmentLock when taking snapshot for the segment
because the segment directory can be
- // modified, e.g. a new snapshot file can be added to the directory. If
not taking the lock, the Helix task
- // thread replacing the segment could fail. For example, we found
FileUtils.cleanDirectory() failed due to
+ }
+ TableDataManager tableDataManager = _context.getTableDataManager();
+ Preconditions.checkNotNull(tableDataManager, "Taking snapshot requires
tableDataManager");
+ boolean isSegmentSkipped = false;
+ for (ImmutableSegmentImpl segment : segmentsWithSnapshot) {
+ // Acquire the segmentLock when taking snapshot for the segment because
the segment directory can be modified,
+ // e.g. a new snapshot file can be added to the directory. If not taking
the lock, the Helix task thread
+ // replacing the segment could fail. For example, we found
FileUtils.cleanDirectory() failed due to
// DirectoryNotEmptyException because a new snapshot file got added into
the segment directory just between two
// major cleanup steps in the cleanDirectory() method.
+ // The lock is acquired in a non-blocking manner so that the consuming
thread is not blocked by the thread
+ // processing the segment, which can hold the segmentLock for long, e.g.
reloading the segment.
String segmentName = segment.getSegmentName();
Lock segmentLock = tableDataManager.getSegmentLock(segmentName);
boolean locked = segmentLock.tryLock();
if (!locked) {
- // Try to get the segmentLock in a non-blocking manner to avoid
deadlock. The Helix task thread takes
- // segmentLock first and then the snapshot RLock when replacing a
segment. However, the consuming thread has
- // already acquired the snapshot WLock when reaching here, and if it
has to wait for segmentLock, it may
- // enter deadlock with the Helix task threads waiting for snapshot
RLock.
- // If we can't get the segmentLock, we'd better skip taking snapshot
for this tracked segment, because its
- // validDocIds/queryableDocIds might become stale or wrong when the
segment is being processed by another
- // thread right now.
+ // The snapshot kept on disk might have become stale, so skip creating
new snapshot files in the next loop to
+ // keep all the snapshots on disk disjoint with each other.
_logger.warn("Could not get segmentLock to take snapshot for segment:
{}, skipping", segmentName);
isSegmentSkipped = true;
continue;
}
try {
- ImmutableSegmentImpl immutableSegment = (ImmutableSegmentImpl) segment;
- if
(!immutableSegment.hasSnapshotFile(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME)
|| (
- _deleteRecordColumn != null && !immutableSegment.hasSnapshotFile(
- V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME))) {
- segmentsWithoutSnapshot.add(immutableSegment);
+ // The segment can be replaced or removed by another thread before the
segmentLock is acquired. The replaced
+ // segment object no longer owns the segment directory, and its
bitmaps might have been drained by the
+ // replacement, so skip persisting them. The new segment object is
tracked as updated, and will be covered by
+ // the next snapshot round. The snapshot kept on disk for the segment
might be stale though, so also skip
+ // creating new snapshot files in the next loop.
+ if (!_trackedSegments.contains(segment)) {
+ _logger.warn("Segment: {} got replaced or removed before taking
snapshot, skipping", segmentName);
+ isSegmentSkipped = true;
continue;
}
- ThreadSafeMutableRoaringBitmap validDocIds =
immutableSegment.getValidDocIds();
+ ThreadSafeMutableRoaringBitmap validDocIds = segment.getValidDocIds();
// NOTE: Segment out of TTL without snapshot might have null
validDocIds
if (validDocIds != null) {
ThreadSafeMutableRoaringBitmap.CardinalityAndBytes
validDocIdsSnapshot = validDocIds.getBytesAndCardinality();
-
immutableSegment.persistDocIdsSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME,
validDocIdsSnapshot);
+
segment.persistDocIdsSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME,
validDocIdsSnapshot);
numPrimaryKeysInSnapshot += validDocIdsSnapshot.getCardinality();
}
if (_deleteRecordColumn != null) {
- ThreadSafeMutableRoaringBitmap queryableDocIds =
immutableSegment.getQueryableDocIds();
+ ThreadSafeMutableRoaringBitmap queryableDocIds =
segment.getQueryableDocIds();
if (queryableDocIds != null) {
ThreadSafeMutableRoaringBitmap.CardinalityAndBytes
queryableDocIdsSnapshot =
queryableDocIds.getBytesAndCardinality();
-
immutableSegment.persistDocIdsSnapshot(V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME,
- queryableDocIdsSnapshot);
+
segment.persistDocIdsSnapshot(V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME,
queryableDocIdsSnapshot);
numQueryableDocIdsInSnapshot +=
queryableDocIdsSnapshot.getCardinality();
}
}
@@ -1006,13 +1016,16 @@ public abstract class
BasePartitionUpsertMetadataManager implements PartitionUps
}
}
// If we have skipped any segments in the previous for-loop, we should
skip the next for-loop, basically to not
- // add new snapshot files on disk. This ensures all the validDocIds &
queryable docIds snapshots kept on disk are
- // still disjoint with each other, although some of them may have become
stale, i.e. tracking more valid docs than
- // expected.
+ // add new snapshot files on disk. This ensures all the validDocIds
snapshots kept on disk are still disjoint
+ // with each other, although some of them may have become stale, i.e.
tracking more valid docs than expected.
if (!isSegmentSkipped) {
for (ImmutableSegmentImpl segment : segmentsWithoutSnapshot) {
String segmentName = segment.getSegmentName();
Lock segmentLock = tableDataManager.getSegmentLock(segmentName);
+ // Unlike the previous loop, skipping a segment here on lock
contention is benign: no snapshot file is
+ // written for it, so the snapshots kept on disk remain disjoint, and
the segment stays tracked as updated
+ // to be retried in the next snapshot round. This is expected for the
just committed segment, whose
+ // segmentLock can still be held by the thread replacing it with the
immutable one.
boolean locked = segmentLock.tryLock();
if (!locked) {
_logger.warn("Could not get segmentLock to take snapshot for
segment: {} w/o snapshot, skipping",
@@ -1020,12 +1033,21 @@ public abstract class
BasePartitionUpsertMetadataManager implements PartitionUps
continue;
}
try {
+ // The segment can be replaced or removed by another thread after it
was classified without snapshot.
+ // The replaced segment object no longer owns the segment directory,
so skip persisting its bitmaps.
+ if (!_trackedSegments.contains(segment)) {
+ _logger.warn("Segment: {} got replaced or removed before taking
snapshot, skipping", segmentName);
+ continue;
+ }
ThreadSafeMutableRoaringBitmap validDocIds =
segment.getValidDocIds();
// NOTE: Segment out of TTL without snapshot might have null
validDocIds
if (validDocIds != null) {
ThreadSafeMutableRoaringBitmap.CardinalityAndBytes
validDocIdsSnapshot =
validDocIds.getBytesAndCardinality();
segment.persistDocIdsSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME,
validDocIdsSnapshot);
+ // The segment has its validDocIds snapshot file on disk now, so
handle it as a segment with snapshot
+ // from now on, even if persisting the queryableDocIds snapshot
below fails.
+ _segmentsWithSnapshot.add(segment);
numPrimaryKeysInSnapshot += validDocIdsSnapshot.getCardinality();
}
if (_deleteRecordColumn != null) {
@@ -1047,6 +1069,7 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
}
}
_updatedSegmentsSinceLastSnapshot.retainAll(_trackedSegments);
+ _segmentsWithSnapshot.retainAll(_trackedSegments);
// Persist TTL watermark after taking snapshots if TTL is enabled, so that
segments out of TTL can be loaded with
// updated validDocIds bitmaps. If the TTL watermark is persisted first,
segments out of TTL may get loaded with
// stale bitmaps or even no bitmap snapshots to use.
@@ -1244,9 +1267,43 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
}
protected void trackUpdatedSegmentsSinceLastSnapshot(IndexSegment segment) {
- if (_enableSnapshot && segment instanceof ImmutableSegment) {
- _updatedSegmentsSinceLastSnapshot.add(segment);
+ if (_enableSnapshot && segment instanceof ImmutableSegmentImpl
immutableSegment) {
+ _updatedSegmentsSinceLastSnapshot.add(immutableSegment);
+ }
+ }
+
+ /// Tracks the segment as updated since the last snapshot, and as a segment
with snapshot when its validDocIds
+ /// snapshot file already exists on disk. No-op when snapshot is not enabled.
+ ///
+ /// The segment is added to [#_segmentsWithSnapshot] before
[#_updatedSegmentsSinceLastSnapshot] because the
+ /// snapshot flow checks the updated set before the with-snapshot set, so
observing the updated marker guarantees
+ /// observing the with-snapshot membership as well, and a segment can never
be misclassified as without snapshot
+ /// while its snapshot file is on disk.
+ protected void trackSegmentForSnapshot(ImmutableSegmentImpl segment) {
+ if (!_enableSnapshot) {
+ return;
+ }
+ if (hasValidDocIdsSnapshotFile(segment)) {
+ _segmentsWithSnapshot.add(segment);
}
+ _updatedSegmentsSinceLastSnapshot.add(segment);
+ }
+
+ /// Removes the segment from all the tracking sets. The segment is removed
eagerly instead of waiting for the next
+ /// snapshot round to clean it up, to not keep stale segment object
references around.
+ protected void untrackSegment(IndexSegment segment) {
+ _trackedSegments.remove(segment);
+ _updatedSegmentsSinceLastSnapshot.remove(segment);
+ _segmentsWithSnapshot.remove(segment);
+ }
+
+ /// Returns `true` when the segment has the validDocIds snapshot file on
disk. The queryableDocIds snapshot file is
+ /// not checked because only the validDocIds snapshots are unioned across
segments during preload, so only their
+ /// staleness can break the disjointness of the snapshots kept on disk. Both
snapshot files are always persisted
+ /// together, so a missing queryableDocIds snapshot file gets recreated the
next time the segment's snapshot is
+ /// taken.
+ protected boolean hasValidDocIdsSnapshotFile(ImmutableSegmentImpl segment) {
+ return
segment.hasSnapshotFile(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME);
}
protected void doClose()
@@ -1302,8 +1359,7 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
*/
protected long getAuthoritativeUpdateOrCreationTime(IndexSegment segment) {
SegmentMetadata segmentMetadata = segment.getSegmentMetadata();
- if (segmentMetadata instanceof SegmentMetadataImpl) {
- SegmentMetadataImpl segmentMetadataImpl = (SegmentMetadataImpl)
segmentMetadata;
+ if (segmentMetadata instanceof SegmentMetadataImpl segmentMetadataImpl) {
if (_tableType == TableType.OFFLINE) {
long zkPushTime = segmentMetadataImpl.getZkPushTime();
if (zkPushTime != Long.MIN_VALUE) {
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManagerTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManagerTest.java
index c448fcb5c94..ca3e2da05a5 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManagerTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManagerTest.java
@@ -126,9 +126,9 @@ public class BasePartitionUpsertMetadataManagerTest {
File segDir02 = new File(TEMP_DIR, "seg02");
ImmutableSegmentImpl seg02 = createImmutableSegment("seg02", segDir02,
segmentsTakenSnapshot, null);
seg02.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3, 4, 5),
null);
- upsertMetadataManager.addSegment(seg02);
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.addSegment(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -187,9 +187,9 @@ public class BasePartitionUpsertMetadataManagerTest {
File segDir02 = new File(TEMP_DIR, "seg02");
ImmutableSegmentImpl seg02 = createImmutableSegment("seg02", segDir02,
segmentsTakenSnapshot, null);
seg02.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3, 4, 5),
null);
- upsertMetadataManager.addSegment(seg02);
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.addSegment(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -239,6 +239,61 @@ public class BasePartitionUpsertMetadataManagerTest {
}
}
+ @Test
+ public void testTakeSnapshotSkipsSegmentWithoutSnapshotOnLockContention()
+ throws IOException {
+ UpsertContext upsertContext = mock(UpsertContext.class);
+ when(upsertContext.isSnapshotEnabled()).thenReturn(true);
+ TableDataManager tdm = mock(TableDataManager.class);
+ when(upsertContext.getTableDataManager()).thenReturn(tdm);
+ Map<String, Lock> segmentLocks = new HashMap<>();
+ // A contended segmentLock on a segment without existing snapshot file
skips only that segment, and the other
+ // segments still take their snapshots.
+ Lock seg01Lock = new ReentrantLock() {
+ @Override
+ public boolean tryLock() {
+ return false;
+ }
+ };
+ segmentLocks.put("seg01", seg01Lock);
+ segmentLocks.put("seg02", new ReentrantLock());
+ segmentLocks.put("seg03", new ReentrantLock());
+ when(tdm.getSegmentLock(anyString())).thenAnswer(invocation -> {
+ String segmentName = invocation.getArgument(0);
+ return segmentLocks.get(segmentName);
+ });
+ DummyPartitionUpsertMetadataManager upsertMetadataManager =
+ new DummyPartitionUpsertMetadataManager("myTable", 0, upsertContext);
+
+ List<String> segmentsTakenSnapshot = new ArrayList<>();
+
+ File segDir01 = new File(TEMP_DIR, "seg01");
+ ImmutableSegmentImpl seg01 = createImmutableSegment("seg01", segDir01,
segmentsTakenSnapshot, null);
+ seg01.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3), null);
+ upsertMetadataManager.addSegment(seg01);
+
+ File segDir02 = new File(TEMP_DIR, "seg02");
+ ImmutableSegmentImpl seg02 = createImmutableSegment("seg02", segDir02,
segmentsTakenSnapshot, null);
+ seg02.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3, 4, 5),
null);
+ // seg02 has snapshot file, so its snapshot is taken first.
+ FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.addSegment(seg02);
+
+ File segDir03 = new File(TEMP_DIR, "seg03");
+ ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
+ seg03.enableUpsert(upsertMetadataManager, createDocIds(3, 4, 7), null);
+ upsertMetadataManager.addSegment(seg03);
+
+ upsertMetadataManager.doTakeSnapshot();
+ // seg01 is skipped on lock contention, and the other segments still take
their snapshots.
+ assertEquals(segmentsTakenSnapshot.size(), 2);
+ assertEquals(segmentsTakenSnapshot.get(0), "seg02");
+ assertTrue(segmentsTakenSnapshot.contains("seg03"));
+ assertFalse(new File(segDir01,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME).exists());
+
assertEquals(seg02.loadDocIdsFromSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME).getCardinality(),
6);
+
assertEquals(seg03.loadDocIdsFromSnapshot(V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME).getCardinality(),
3);
+ }
+
@Test
public void testTakeSnapshotInOrderBasedOnUpdates()
throws IOException {
@@ -262,9 +317,9 @@ public class BasePartitionUpsertMetadataManagerTest {
File segDir02 = new File(TEMP_DIR, "seg02");
ImmutableSegmentImpl seg02 = createImmutableSegment("seg02", segDir02,
segmentsTakenSnapshot, null);
seg02.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3, 4, 5),
null);
- upsertMetadataManager.addSegment(seg02);
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.addSegment(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -319,9 +374,9 @@ public class BasePartitionUpsertMetadataManagerTest {
File segDir02 = new File(TEMP_DIR, "seg02");
ImmutableSegmentImpl seg02 = createImmutableSegment("seg02", segDir02,
segmentsTakenSnapshot, null);
seg02.enableUpsert(upsertMetadataManager, createDocIds(0, 1, 2, 3, 4, 5),
null);
- upsertMetadataManager.addSegment(seg02);
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.addSegment(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -385,6 +440,7 @@ public class BasePartitionUpsertMetadataManagerTest {
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
FileUtils.touch(new File(segDir02,
V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.markSegmentWithSnapshot(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 =
@@ -459,6 +515,7 @@ public class BasePartitionUpsertMetadataManagerTest {
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
FileUtils.touch(new File(segDir02,
V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.markSegmentWithSnapshot(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -543,6 +600,7 @@ public class BasePartitionUpsertMetadataManagerTest {
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
FileUtils.touch(new File(segDir02,
V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.markSegmentWithSnapshot(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -606,6 +664,7 @@ public class BasePartitionUpsertMetadataManagerTest {
// seg02 has snapshot file, so its snapshot is taken first.
FileUtils.touch(new File(segDir02,
V1Constants.VALID_DOC_IDS_SNAPSHOT_FILE_NAME));
FileUtils.touch(new File(segDir02,
V1Constants.QUERYABLE_DOC_IDS_SNAPSHOT_FILE_NAME));
+ upsertMetadataManager.markSegmentWithSnapshot(seg02);
File segDir03 = new File(TEMP_DIR, "seg03");
ImmutableSegmentImpl seg03 = createImmutableSegment("seg03", segDir03,
segmentsTakenSnapshot, null);
@@ -1026,10 +1085,14 @@ public class BasePartitionUpsertMetadataManagerTest {
_trackedSegments.add(seg);
}
- public void markSegmentAsUpdated(IndexSegment seg) {
+ public void markSegmentAsUpdated(ImmutableSegmentImpl seg) {
_updatedSegmentsSinceLastSnapshot.add(seg);
}
+ public void markSegmentWithSnapshot(ImmutableSegmentImpl seg) {
+ _segmentsWithSnapshot.add(seg);
+ }
+
@Override
protected long getNumPrimaryKeys() {
return 0;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]