Skip to content

Commit de2f596

Browse files
authored
Add skipUpsertDelete query option to view valid docs instead of queryable docs (#19009)
1 parent 41ee76b commit de2f596

18 files changed

Lines changed: 514 additions & 63 deletions

File tree

pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,10 @@ public static boolean isSkipUpsertView(Map<String, String> queryOptions) {
165165
return Boolean.parseBoolean(queryOptions.get(QueryOptionKey.SKIP_UPSERT_VIEW));
166166
}
167167

168+
public static boolean isSkipUpsertDelete(Map<String, String> queryOptions) {
169+
return Boolean.parseBoolean(queryOptions.get(QueryOptionKey.SKIP_UPSERT_DELETE));
170+
}
171+
168172
public static boolean isTraceRuleProductions(Map<String, String> queryOptions) {
169173
return Boolean.parseBoolean(queryOptions.get(QueryOptionKey.TRACE_RULE_PRODUCTIONS));
170174
}

pinot-core/src/main/java/org/apache/pinot/core/plan/FilterPlanNode.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -97,20 +97,20 @@ public FilterPlanNode(SegmentContext segmentContext, QueryContext queryContext,
9797

9898
@Override
9999
public BaseFilterOperator run() {
100-
MutableRoaringBitmap queryableDocIdsSnapshot = _segmentContext.getQueryableDocIdsSnapshot();
100+
MutableRoaringBitmap docIdsSnapshot = _segmentContext.getDocIdsSnapshot();
101101
int numDocs = _indexSegment.getSegmentMetadata().getTotalDocs();
102102

103103
if (_filter != null) {
104104
BaseFilterOperator filterOperator = constructPhysicalOperator(_filter, numDocs);
105-
if (queryableDocIdsSnapshot != null) {
106-
BaseFilterOperator validDocFilter = new BitmapBasedFilterOperator(queryableDocIdsSnapshot, false, numDocs);
105+
if (docIdsSnapshot != null) {
106+
BaseFilterOperator validDocFilter = new BitmapBasedFilterOperator(docIdsSnapshot, false, numDocs);
107107
return FilterOperatorUtils.getAndFilterOperator(_queryContext, Arrays.asList(filterOperator, validDocFilter),
108108
numDocs);
109109
} else {
110110
return filterOperator;
111111
}
112-
} else if (queryableDocIdsSnapshot != null) {
113-
return new BitmapBasedFilterOperator(queryableDocIdsSnapshot, false, numDocs);
112+
} else if (docIdsSnapshot != null) {
113+
return new BitmapBasedFilterOperator(docIdsSnapshot, false, numDocs);
114114
} else {
115115
return new MatchAllFilterOperator(numDocs);
116116
}

pinot-core/src/main/java/org/apache/pinot/core/query/pruner/SegmentPrunerService.java

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -136,26 +136,31 @@ public List<IndexSegment> prune(List<IndexSegment> segments, QueryContext query,
136136
* @param segments the list of segments to be pruned. This is a destructive operation that may modify this list in an
137137
* undefined way. Therefore, this list should not be used after calling this method.
138138
* @param query query context; when non-null and skipUpsert=true, segments with 0 queryable/valid docs are not
139-
* treated as empty (they contribute replaced rows to the result). When queryable doc ids exist,
140-
* emptiness is determined from them; otherwise {@link IndexSegment#getValidDocIds()} is used.
139+
* treated as empty (they contribute replaced rows to the result). When skipUpsertDelete=true,
140+
* emptiness is determined from valid docs (tombstones count as non-empty); otherwise from
141+
* queryable docs.
141142
* @return the new list with filtered elements. This is the list that have to be used.
142143
*/
143144
private static List<IndexSegment> removeEmptySegments(List<IndexSegment> segments, QueryContext query) {
144145
int selected = 0;
145-
boolean skipUpsert = QueryOptionsUtils.isSkipUpsert(query.getQueryOptions());
146+
Map<String, String> queryOptions = query.getQueryOptions();
147+
boolean skipUpsert = QueryOptionsUtils.isSkipUpsert(queryOptions);
148+
boolean skipUpsertDelete = QueryOptionsUtils.isSkipUpsertDelete(queryOptions);
146149
for (IndexSegment segment : segments) {
147-
if (!isEmptySegment(segment, skipUpsert)) {
150+
if (!isEmptySegment(segment, skipUpsert, skipUpsertDelete)) {
148151
segments.set(selected++, segment);
149152
}
150153
}
151154
return segments.subList(0, selected);
152155
}
153156

154-
private static boolean isEmptySegment(IndexSegment segment, boolean skipUpsert) {
157+
private static boolean isEmptySegment(IndexSegment segment, boolean skipUpsert, boolean skipUpsertDelete) {
155158
if (segment.getSegmentMetadata().getTotalDocs() == 0) {
156159
return true;
157160
}
158-
// Check if the segment has 0 queryable docIds while skipUpsert=false
159-
return !skipUpsert && segment.hasNoQueryableDocs();
161+
if (skipUpsert) {
162+
return false;
163+
}
164+
return skipUpsertDelete ? segment.hasNoValidDocs() : segment.hasNoQueryableDocs();
160165
}
161166
}

pinot-core/src/test/java/org/apache/pinot/core/plan/FilterPlanNodeTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,7 @@ public void testConsistentSnapshot()
9494
// Result should be invariant - always exactly 3 docs
9595
for (int i = 0; i < 10_000; i++) {
9696
SegmentContext segmentContext = new SegmentContext(segment);
97-
segmentContext.setQueryableDocIdsSnapshot(UpsertUtils.getQueryableDocIdsSnapshotFromSegment(segment));
97+
segmentContext.setDocIdsSnapshot(UpsertUtils.getQueryableDocIdsSnapshotFromSegment(segment));
9898
assertEquals(getNumberOfFilteredDocs(segmentContext, queryContext), 3);
9999
}
100100

pinot-core/src/test/java/org/apache/pinot/core/plan/maker/MetadataAndDictionaryAggregationPlanMakerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -165,7 +165,7 @@ public void testPlanMaker(String query, Class<? extends Operator<?>> operatorCla
165165
assertTrue(operatorClass.isInstance(operator));
166166

167167
SegmentContext segmentContext = new SegmentContext(_upsertIndexSegment);
168-
segmentContext.setQueryableDocIdsSnapshot(UpsertUtils.getQueryableDocIdsSnapshotFromSegment(_upsertIndexSegment));
168+
segmentContext.setDocIdsSnapshot(UpsertUtils.getQueryableDocIdsSnapshotFromSegment(_upsertIndexSegment));
169169
Operator<?> upsertOperator = PLAN_MAKER.makeSegmentPlanNode(segmentContext, queryContext).run();
170170
assertTrue(upsertOperatorClass.isInstance(upsertOperator));
171171
}

pinot-core/src/test/java/org/apache/pinot/core/query/pruner/SegmentPrunerServiceTest.java

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -186,6 +186,42 @@ public void emptyQueryableRetainedWithSkipUpsert() {
186186
Assert.assertEquals(actual, segments);
187187
}
188188

189+
@Test
190+
public void emptyQueryableRetainedWithSkipUpsertDelete() {
191+
SegmentPrunerService service = new SegmentPrunerService(_emptyPrunerConf);
192+
ThreadSafeMutableRoaringBitmap valid = new ThreadSafeMutableRoaringBitmap(0);
193+
ThreadSafeMutableRoaringBitmap queryable = new ThreadSafeMutableRoaringBitmap();
194+
IndexSegment segment = mockUpsertIndexSegment(10, valid, queryable);
195+
196+
List<IndexSegment> segments = new ArrayList<>();
197+
segments.add(segment);
198+
QueryContext queryContext =
199+
QueryContextConverterUtils.getQueryContext("select col1 from t1 option(skipUpsertDelete=true)");
200+
201+
List<IndexSegment> actual = service.prune(segments, queryContext, new SegmentPrunerStatistics());
202+
203+
Assert.assertEquals(actual, segments);
204+
}
205+
206+
/**
207+
* skipUpsertDelete checks valid-docs emptiness specifically, not skipUpsert's blanket bypass: a segment fully
208+
* superseded elsewhere (0 valid docs) is still pruned.
209+
*/
210+
@Test
211+
public void emptyValidPrunedWithSkipUpsertDelete() {
212+
SegmentPrunerService service = new SegmentPrunerService(_emptyPrunerConf);
213+
IndexSegment segment = mockUpsertIndexSegment(10, new ThreadSafeMutableRoaringBitmap(), null);
214+
215+
List<IndexSegment> segments = new ArrayList<>();
216+
segments.add(segment);
217+
QueryContext queryContext =
218+
QueryContextConverterUtils.getQueryContext("select col1 from t1 option(skipUpsertDelete=true)");
219+
220+
List<IndexSegment> actual = service.prune(segments, queryContext, new SegmentPrunerStatistics());
221+
222+
Assert.assertEquals(actual, List.of());
223+
}
224+
189225
@Test
190226
public void nonEmptyQueryableNotPruned() {
191227
SegmentPrunerService service = new SegmentPrunerService(_emptyPrunerConf);

pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/immutable/ImmutableSegmentImpl.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@
4040
import org.apache.pinot.segment.local.segment.virtualcolumn.VirtualColumnContext;
4141
import org.apache.pinot.segment.local.startree.v2.store.StarTreeIndexContainer;
4242
import org.apache.pinot.segment.local.upsert.PartitionUpsertMetadataManager;
43+
import org.apache.pinot.segment.local.upsert.UpsertUtils;
4344
import org.apache.pinot.segment.local.upsert.UpsertViewManager;
4445
import org.apache.pinot.segment.spi.ColumnMetadata;
4546
import org.apache.pinot.segment.spi.FetchContext;
@@ -382,6 +383,11 @@ public boolean hasNoQueryableDocs() {
382383
return validDocIds != null && validDocIds.isEmpty();
383384
}
384385

386+
@Override
387+
public boolean hasNoValidDocs() {
388+
return UpsertUtils.hasNoValidDocs(_partitionUpsertMetadataManager, this);
389+
}
390+
385391
@Override
386392
public GenericRow getRecord(int docId, GenericRow reuse) {
387393
try (PinotSegmentRecordReader recordReader = new PinotSegmentRecordReader()) {

pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImpl.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@
7474
import org.apache.pinot.segment.local.upsert.PartitionUpsertMetadataManager;
7575
import org.apache.pinot.segment.local.upsert.RecordInfo;
7676
import org.apache.pinot.segment.local.upsert.UpsertContext;
77+
import org.apache.pinot.segment.local.upsert.UpsertUtils;
7778
import org.apache.pinot.segment.local.upsert.UpsertViewManager;
7879
import org.apache.pinot.segment.local.utils.FixedIntArrayOffHeapIdMap;
7980
import org.apache.pinot.segment.local.utils.IdMap;
@@ -1242,6 +1243,11 @@ public boolean hasNoQueryableDocs() {
12421243
return validDocIds != null && validDocIds.isEmpty();
12431244
}
12441245

1246+
@Override
1247+
public boolean hasNoValidDocs() {
1248+
return UpsertUtils.hasNoValidDocs(_partitionUpsertMetadataManager, this);
1249+
}
1250+
12451251
@Override
12461252
public GenericRow getRecord(int docId, GenericRow reuse) {
12471253
try (PinotSegmentRecordReader recordReader = new PinotSegmentRecordReader()) {

pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/ConcurrentMapTableUpsertMetadataManager.java

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -95,24 +95,28 @@ public void setSegmentContexts(List<SegmentContext> segmentContexts, Map<String,
9595
// Otherwise, get queryableDocIds bitmaps as kept by the segment objects directly as before.
9696
if (_context.getConsistencyMode() == UpsertConfig.ConsistencyMode.NONE || QueryOptionsUtils.isSkipUpsertView(
9797
queryOptions)) {
98+
// No shared upsert-view lock exists in this branch, so a direct read is already safe here.
99+
boolean skipUpsertDelete = QueryOptionsUtils.isSkipUpsertDelete(queryOptions);
98100
for (SegmentContext segmentContext : segmentContexts) {
99101
IndexSegment segment = segmentContext.getIndexSegment();
100-
segmentContext.setQueryableDocIdsSnapshot(UpsertUtils.getQueryableDocIdsSnapshotFromSegment(segment));
102+
segmentContext.setDocIdsSnapshot(skipUpsertDelete
103+
? UpsertUtils.getValidDocIdsSnapshotFromSegment(segment)
104+
: UpsertUtils.getQueryableDocIdsSnapshotFromSegment(segment));
101105
}
102106
return;
103107
}
104-
// All segments should have been tracked by partitionMetadataManagers to provide queries consistent upsert view.
108+
// A consistency mode is active: UpsertViewManager knows the locking each mode requires for skipUpsertDelete too.
105109
_partitionMetadataManagerMap.forEach(
106110
(partitionID, upsertMetadataManager) -> upsertMetadataManager.getUpsertViewManager()
107111
.setSegmentContexts(segmentContexts, queryOptions));
108112
if (LOGGER.isDebugEnabled()) {
109113
for (SegmentContext segmentContext : segmentContexts) {
110114
IndexSegment segment = segmentContext.getIndexSegment();
111-
if (segmentContext.getQueryableDocIdsSnapshot() == null) {
115+
if (segmentContext.getDocIdsSnapshot() == null) {
112116
LOGGER.debug("No upsert view for segment: {}, type: {}, total: {}", segment.getSegmentName(),
113117
(segment instanceof ImmutableSegment ? "imm" : "mut"), segment.getSegmentMetadata().getTotalDocs());
114118
} else {
115-
int cardCnt = segmentContext.getQueryableDocIdsSnapshot().getCardinality();
119+
int cardCnt = segmentContext.getDocIdsSnapshot().getCardinality();
116120
LOGGER.debug("Got upsert view of segment: {}, type: {}, total: {}, valid: {}", segment.getSegmentName(),
117121
(segment instanceof ImmutableSegment ? "imm" : "mut"), segment.getSegmentMetadata().getTotalDocs(),
118122
cardCnt);

pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertUtils.java

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,36 @@ public static MutableRoaringBitmap getQueryableDocIdsSnapshotFromSegment(IndexSe
5757
: (useEmptyForNull ? new MutableRoaringBitmap() : null);
5858
}
5959

60+
/// Unlike [#getQueryableDocIdsSnapshotFromSegment], never falls back to preferring queryable docs.
61+
@Nullable
62+
public static MutableRoaringBitmap getValidDocIdsSnapshotFromSegment(IndexSegment segment) {
63+
return getValidDocIdsSnapshotFromSegment(segment, false);
64+
}
65+
66+
/// Shared by {@code ImmutableSegmentImpl}/{@code MutableSegmentImpl}'s `hasNoValidDocs()`. Mirrors
67+
/// `hasNoQueryableDocs()`'s consistency-mode-aware/live-fallback split, against the valid-docs cache instead.
68+
public static boolean hasNoValidDocs(@Nullable PartitionUpsertMetadataManager partitionUpsertMetadataManager,
69+
IndexSegment segment) {
70+
if (partitionUpsertMetadataManager == null) {
71+
return false;
72+
}
73+
UpsertViewManager viewManager = partitionUpsertMetadataManager.getUpsertViewManager();
74+
if (viewManager != null) {
75+
MutableRoaringBitmap validDocIdsSnapshot = viewManager.getValidDocIdsSnapshot(segment);
76+
return validDocIdsSnapshot != null && validDocIdsSnapshot.isEmpty();
77+
}
78+
ThreadSafeMutableRoaringBitmap validDocIds = segment.getValidDocIds();
79+
return validDocIds != null && validDocIds.isEmpty();
80+
}
81+
82+
@Nullable
83+
public static MutableRoaringBitmap getValidDocIdsSnapshotFromSegment(IndexSegment segment,
84+
boolean useEmptyForNull) {
85+
ThreadSafeMutableRoaringBitmap validDocIds = segment.getValidDocIds();
86+
return validDocIds != null ? validDocIds.getMutableRoaringBitmap()
87+
: (useEmptyForNull ? new MutableRoaringBitmap() : null);
88+
}
89+
6090
public static void doReplaceDocId(ThreadSafeMutableRoaringBitmap validDocIds,
6191
@Nullable ThreadSafeMutableRoaringBitmap queryableDocIds, int oldDocId, int newDocId, RecordInfo recordInfo) {
6292
validDocIds.replace(oldDocId, newDocId);

0 commit comments

Comments
 (0)