Jackie-Jiang commented on code in PR #19430:
URL: https://github.com/apache/pinot/pull/19430#discussion_r3984992107
##########
pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java:
##########
@@ -53,4 +96,207 @@ public void testVariableSizeRecordOffsets() {
assertEquals(offsets.getStartOffset(3), Integer.MAX_VALUE + 456L);
assertEquals(offsets.getEndOffset(), Integer.MAX_VALUE + 456L + 789L);
}
+
+ /// Builds a star-tree with the off-heap builder and asserts the
intermediate segment record file
+ /// is cleaned up after build. Exercises the full
sortAndAggregateSegmentRecords → iterator → close
+ /// lifecycle, including buffer allocation, sequential fill in Sub-phase A,
sort in Sub-phase B,
+ /// and dim-from-buffer reads in Sub-phase C.
+ @Test
+ public void testBuildCleansUpSegmentRecordFile()
+ throws Exception {
+ buildTestSegment();
+
+ List<StarTreeV2BuilderConfig> builderConfigs = createBuilderConfigs();
+ File segmentDir = INDEX_DIR.listFiles()[0];
+
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(builderConfigs, segmentDir,
+ MultipleTreesBuilder.BuildMode.OFF_HEAP)) {
+ builder.build();
+ }
+
+ // OffHeapSingleTreeBuilder writes an intermediate segment.record file
during Sub-phase A/B/C.
+ // The dim-buffer optimization keeps that file alive through iterator
exhaustion; close() must
+ // release the buffer and delete the file.
+ File segmentRecordFile = findSegmentRecordFile(segmentDir);
+ assertFalse(segmentRecordFile != null && segmentRecordFile.exists(),
+ "segment.record file should be deleted after build: " +
segmentRecordFile);
Review Comment:
[P2] Verify buffer release directly
The fixture allocates a 40-byte DIRECT buffer, so `segment.record` is never
created and this assertion cannot detect a leak. I confirmed that all five
tests still pass with `releaseSegmentRecordBuffer()` completely disabled. The
close-without-build test also never constructs `OffHeapSingleTreeBuilder`.
Please add direct builder tests that assert buffer count/usage returns to
baseline after iterator exhaustion and after closing an undrained iterator.
##########
pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java:
##########
@@ -53,4 +96,207 @@ public void testVariableSizeRecordOffsets() {
assertEquals(offsets.getStartOffset(3), Integer.MAX_VALUE + 456L);
assertEquals(offsets.getEndOffset(), Integer.MAX_VALUE + 456L + 789L);
}
+
+ /// Builds a star-tree with the off-heap builder and asserts the
intermediate segment record file
+ /// is cleaned up after build. Exercises the full
sortAndAggregateSegmentRecords → iterator → close
+ /// lifecycle, including buffer allocation, sequential fill in Sub-phase A,
sort in Sub-phase B,
+ /// and dim-from-buffer reads in Sub-phase C.
+ @Test
+ public void testBuildCleansUpSegmentRecordFile()
+ throws Exception {
+ buildTestSegment();
+
+ List<StarTreeV2BuilderConfig> builderConfigs = createBuilderConfigs();
+ File segmentDir = INDEX_DIR.listFiles()[0];
+
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(builderConfigs, segmentDir,
+ MultipleTreesBuilder.BuildMode.OFF_HEAP)) {
+ builder.build();
+ }
+
+ // OffHeapSingleTreeBuilder writes an intermediate segment.record file
during Sub-phase A/B/C.
+ // The dim-buffer optimization keeps that file alive through iterator
exhaustion; close() must
+ // release the buffer and delete the file.
+ File segmentRecordFile = findSegmentRecordFile(segmentDir);
+ assertFalse(segmentRecordFile != null && segmentRecordFile.exists(),
+ "segment.record file should be deleted after build: " +
segmentRecordFile);
+ }
+
+ /// close() before build() must not throw and must leave no leftover segment
record file.
+ @Test
+ public void testCloseWithoutBuildDoesNotThrow()
+ throws Exception {
+ buildTestSegment();
+
+ List<StarTreeV2BuilderConfig> builderConfigs = createBuilderConfigs();
+ File segmentDir = INDEX_DIR.listFiles()[0];
+
+ // Construct and immediately close — no build().
+ MultipleTreesBuilder builder = new MultipleTreesBuilder(builderConfigs,
segmentDir,
+ MultipleTreesBuilder.BuildMode.OFF_HEAP);
+ builder.close();
+
+ File segmentRecordFile = findSegmentRecordFile(segmentDir);
+ assertFalse(segmentRecordFile != null && segmentRecordFile.exists(),
+ "segment.record file should not exist when build() was never called: "
+ segmentRecordFile);
+ }
+
+ /// Builds the same segment twice — once with OFF_HEAP and once with ON_HEAP
— then compares the
+ /// resulting star-trees by grouping star-tree records by dim tuple and
summing the aggregate
+ /// column. Directly regressions-tests dim-read parity: if the buffer read
returned wrong bytes,
+ /// OFF_HEAP would aggregate metrics under the wrong dim tuples and the maps
would diverge.
+ @Test
+ public void testOffHeapProducesSameStarTreeAsOnHeap()
+ throws Exception {
+ buildTestSegment();
+ File sourceSegmentDir = INDEX_DIR.listFiles()[0];
+ List<StarTreeV2BuilderConfig> builderConfigs = createBuilderConfigs();
+
+ File offHeapDir = new File(TEMP_DIR, "offHeapCopy");
+ File onHeapDir = new File(TEMP_DIR, "onHeapCopy");
+ FileUtils.copyDirectory(sourceSegmentDir, offHeapDir);
+ FileUtils.copyDirectory(sourceSegmentDir, onHeapDir);
+
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(builderConfigs, offHeapDir,
+ MultipleTreesBuilder.BuildMode.OFF_HEAP)) {
+ builder.build();
+ }
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(builderConfigs, onHeapDir,
+ MultipleTreesBuilder.BuildMode.ON_HEAP)) {
+ builder.build();
+ }
+
+ Map<String, Long> offHeapByDim = readStarTreeAggregatedByDim(offHeapDir);
+ Map<String, Long> onHeapByDim = readStarTreeAggregatedByDim(onHeapDir);
+ assertEquals(offHeapByDim, onHeapByDim,
+ "OFF_HEAP and ON_HEAP star-trees disagree on aggregate-by-dim: "
+ + "off=" + offHeapByDim + ", on=" + onHeapByDim);
+ }
+
+ /// Iterates every star-tree record, resolves dim dict-ids to values, and
returns a
+ /// `(stringCol|intCol) -> sum(longCol)` map. Star nodes (dim value = STAR)
fold into their own
+ /// bucket, so the map is a canonical view of the tree independent of
internal doc order.
+ private Map<String, Long> readStarTreeAggregatedByDim(File segmentDir)
+ throws Exception {
+ ImmutableSegment segment = ImmutableSegmentLoader.load(segmentDir,
ReadMode.mmap);
+ try {
+ StarTreeV2 tree = segment.getStarTrees().get(0);
+ int numDocs = tree.getMetadata().getNumDocs();
+
+ ForwardIndexReader<ForwardIndexReaderContext> stringReader =
+ (ForwardIndexReader<ForwardIndexReaderContext>)
tree.getDataSource("stringCol").getForwardIndex();
+ Dictionary stringDict = tree.getDataSource("stringCol").getDictionary();
+ ForwardIndexReader<ForwardIndexReaderContext> intReader =
+ (ForwardIndexReader<ForwardIndexReaderContext>)
tree.getDataSource("intCol").getForwardIndex();
+ Dictionary intDict = tree.getDataSource("intCol").getDictionary();
+ // Aggregate columns are keyed by
AggregationFunctionColumnPair.toColumnName() ("sum__longCol").
+ ForwardIndexReader<ForwardIndexReaderContext> aggReader =
+ (ForwardIndexReader<ForwardIndexReaderContext>)
tree.getDataSource("sum__longCol").getForwardIndex();
+
+ Map<String, Long> byDim = new TreeMap<>();
+ try (ForwardIndexReaderContext stringCtx = stringReader.createContext();
+ ForwardIndexReaderContext intCtx = intReader.createContext();
+ ForwardIndexReaderContext aggCtx = aggReader.createContext()) {
+ for (int docId = 0; docId < numDocs; docId++) {
+ int stringDictId = stringReader.getDictId(docId, stringCtx);
+ int intDictId = intReader.getDictId(docId, intCtx);
+ String stringVal = stringDictId == -1 ? "*" :
stringDict.getStringValue(stringDictId);
+ String intVal = intDictId == -1 ? "*" :
String.valueOf(intDict.getIntValue(intDictId));
+ byDim.merge(stringVal + "|" + intVal, aggReader.getLong(docId,
aggCtx), Long::sum);
+ }
+ }
+ return byDim;
+ } finally {
+ segment.destroy();
+ }
+ }
+
+ private void buildTestSegment()
+ throws Exception {
+ Schema schema = new Schema.SchemaBuilder()
+ .addSingleValueDimension("stringCol", FieldSpec.DataType.STRING)
+ .addSingleValueDimension("intCol", FieldSpec.DataType.INT)
+ .addMetric("longCol", FieldSpec.DataType.LONG)
Review Comment:
[P3] Import `DataType` directly
Please import `FieldSpec.DataType` and use `DataType.STRING`,
`DataType.INT`, and `DataType.LONG`, following the repository convention.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]