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 86107957b03 Star-tree offheap build: read dims from sort buffer
instead of re-fetching from segment (#19430)
86107957b03 is described below
commit 86107957b03ff48817895a34c3640c8baeb1bdd0
Author: Chaitanya Deepthi <[email protected]>
AuthorDate: Thu Sep 10 20:48:41 2026 -0700
Star-tree offheap build: read dims from sort buffer instead of re-fetching
from segment (#19430)
---
.../startree/v2/builder/BaseSingleTreeBuilder.java | 20 ++
.../v2/builder/OffHeapSingleTreeBuilder.java | 108 ++++++---
.../v2/builder/OffHeapSingleTreeBuilderTest.java | 254 +++++++++++++++++++++
3 files changed, 345 insertions(+), 37 deletions(-)
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/BaseSingleTreeBuilder.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/BaseSingleTreeBuilder.java
index 456e8f32e5f..6fe1bfc1548 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/BaseSingleTreeBuilder.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/BaseSingleTreeBuilder.java
@@ -48,6 +48,7 @@ import
org.apache.pinot.segment.spi.index.startree.AggregationFunctionColumnPair
import org.apache.pinot.segment.spi.index.startree.AggregationSpec;
import org.apache.pinot.segment.spi.index.startree.StarTreeNode;
import org.apache.pinot.segment.spi.index.startree.StarTreeV2Constants;
+import org.apache.pinot.segment.spi.memory.PinotDataBuffer;
import org.apache.pinot.spi.data.FieldSpec.DataType;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -231,6 +232,25 @@ abstract class BaseSingleTreeBuilder implements
SingleTreeBuilder {
return new Record(dimensions, metrics);
}
+ /// Same as {@link #getSegmentRecord(int)} but reads dims from `dimBuffer`
(row-major, indexed by
+ /// `docId * _numDimensions * Integer.BYTES`) instead of re-fetching from
the segment.
+ Record getSegmentRecordWithBufferDims(int docId, PinotDataBuffer dimBuffer) {
+ int[] dimensions = new int[_numDimensions];
+ long off = (long) docId * _numDimensions * Integer.BYTES;
+ for (int i = 0; i < _numDimensions; i++) {
+ dimensions[i] = dimBuffer.getInt(off);
+ off += Integer.BYTES;
+ }
+ Object[] metrics = new Object[_numMetrics];
+ for (int i = 0; i < _numMetrics; i++) {
+ // Ignore the column for COUNT aggregation function
+ if (_metricReaders[i] != null) {
+ metrics[i] = _metricReaders[i].getValue(docId);
+ }
+ }
+ return new Record(dimensions, metrics);
+ }
+
/// Merges a segment record (raw) into the aggregated record.
///
/// Will create a new aggregated record if the current one is `null`.
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilder.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilder.java
index f6ab03fa580..f00b7c03823 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilder.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilder.java
@@ -35,10 +35,13 @@ import org.apache.commons.io.FileUtils;
import org.apache.pinot.segment.spi.ImmutableSegment;
import org.apache.pinot.segment.spi.index.startree.StarTreeV2Constants;
import org.apache.pinot.segment.spi.memory.PinotDataBuffer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
/// The `OffHeapSingleTreeBuilder` class is the single star-tree builder that
uses off-heap memory.
public class OffHeapSingleTreeBuilder extends BaseSingleTreeBuilder {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(OffHeapSingleTreeBuilder.class);
private static final String SEGMENT_RECORD_FILE_NAME = "segment.record";
private static final String STAR_TREE_RECORD_FILE_NAME = "star-tree.record";
// If the temporary buffer needed is larger than 500M, use MMAP, otherwise
use DIRECT
@@ -49,6 +52,7 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
private final BufferedOutputStream _starTreeRecordOutputStream;
private final RecordOffsets _starTreeRecordOffsets;
+ private PinotDataBuffer _segmentRecordBuffer;
private PinotDataBuffer _starTreeRecordBuffer;
private int _numReadableStarTreeRecords;
@@ -203,13 +207,12 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
Iterator<Record> sortAndAggregateSegmentRecords(int numDocs)
throws IOException {
// Write all dimensions for segment records into the buffer, and sort all
records using an int array
- PinotDataBuffer dataBuffer;
long bufferSize = (long) numDocs * _numDimensions * Integer.BYTES;
if (bufferSize > MMAP_SIZE_THRESHOLD) {
- dataBuffer = PinotDataBuffer.mapFile(_segmentRecordFile, false, 0,
bufferSize, PinotDataBuffer.NATIVE_ORDER,
- "OffHeapSingleTreeBuilder: segment record buffer");
+ _segmentRecordBuffer = PinotDataBuffer.mapFile(_segmentRecordFile,
false, 0, bufferSize,
+ PinotDataBuffer.NATIVE_ORDER, "OffHeapSingleTreeBuilder: segment
record buffer");
} else {
- dataBuffer = PinotDataBuffer.allocateDirect(bufferSize,
PinotDataBuffer.NATIVE_ORDER,
+ _segmentRecordBuffer = PinotDataBuffer.allocateDirect(bufferSize,
PinotDataBuffer.NATIVE_ORDER,
"OffHeapSingleTreeBuilder: segment record buffer");
}
int[] sortedDocIds = new int[numDocs];
@@ -221,7 +224,7 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
for (int i = 0; i < numDocs; i++) {
int[] dimensions = getSegmentRecordDimensions(i);
for (int j = 0; j < _numDimensions; j++) {
- dataBuffer.putInt(offset, dimensions[j]);
+ _segmentRecordBuffer.putInt(offset, dimensions[j]);
offset += Integer.BYTES;
}
}
@@ -229,8 +232,8 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
long offset1 = (long) sortedDocIds[i1] * _numDimensions *
Integer.BYTES;
long offset2 = (long) sortedDocIds[i2] * _numDimensions *
Integer.BYTES;
for (int i = 0; i < _numDimensions; i++) {
- int dimension1 = dataBuffer.getInt(offset1 + (long) i *
Integer.BYTES);
- int dimension2 = dataBuffer.getInt(offset2 + (long) i *
Integer.BYTES);
+ int dimension1 = _segmentRecordBuffer.getInt(offset1 + (long) i *
Integer.BYTES);
+ int dimension2 = _segmentRecordBuffer.getInt(offset2 + (long) i *
Integer.BYTES);
if (dimension1 != dimension2) {
return dimension1 - dimension2;
}
@@ -241,40 +244,67 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
sortedDocIds[i1] = sortedDocIds[i2];
sortedDocIds[i2] = temp;
});
- } finally {
- dataBuffer.close();
- if (_segmentRecordFile.exists()) {
- FileUtils.forceDelete(_segmentRecordFile);
- }
- }
- // Create an iterator for aggregated records
- return new Iterator<Record>() {
- boolean _hasNext = true;
- Record _currentRecord = getSegmentRecord(sortedDocIds[0]);
- int _docId = 1;
-
- @Override
- public boolean hasNext() {
- return _hasNext;
- }
+ // Iterator reads dims from `_segmentRecordBuffer` (populated
sequentially above) instead of
+ // re-fetching them from the segment forward index in the sorted
(random-with-respect-to-layout)
+ // order. The buffer is a field, guaranteed to be released by close()
via the try-with-resources
+ // on the builder in MultipleTreesBuilder — the terminal next() below is
an early-release
+ // optimization; releaseSegmentRecordBuffer() is idempotent, so
double-release is safe.
+ return new Iterator<Record>() {
+ boolean _hasNext = true;
+ Record _currentRecord =
getSegmentRecordWithBufferDims(sortedDocIds[0], _segmentRecordBuffer);
+ int _docId = 1;
+
+ @Override
+ public boolean hasNext() {
+ return _hasNext;
+ }
- @Override
- public Record next() {
- Record next = mergeSegmentRecord(null, _currentRecord);
- while (_docId < numDocs) {
- Record record = getSegmentRecord(sortedDocIds[_docId++]);
- if (!Arrays.equals(record._dimensions, next._dimensions)) {
- _currentRecord = record;
- return next;
- } else {
- next = mergeSegmentRecord(next, record);
+ @Override
+ public Record next() {
+ Record next = mergeSegmentRecord(null, _currentRecord);
+ while (_docId < numDocs) {
+ Record record =
getSegmentRecordWithBufferDims(sortedDocIds[_docId++], _segmentRecordBuffer);
+ if (!Arrays.equals(record._dimensions, next._dimensions)) {
+ _currentRecord = record;
+ return next;
+ } else {
+ next = mergeSegmentRecord(next, record);
+ }
}
+ _hasNext = false;
+ releaseSegmentRecordBuffer();
+ return next;
}
- _hasNext = false;
- return next;
+ };
+ } catch (Throwable t) {
+ // Fill / sort / iterator construction failed: nobody will drain the
iterator, so release
+ // the buffer here. close() will find `_segmentRecordBuffer == null` and
skip re-releasing.
+ releaseSegmentRecordBuffer();
+ throw t;
+ }
+ }
+
+ // Idempotent: safe to call multiple times from any exit path (iterator
drain, exception during
+ // sortAndAggregate, or close() via try-with-resources). Null-checks the
field and sets it to null
+ // after release, so a second call is a no-op. Best-effort — close()
failures are logged, not
+ // thrown, so a successful build isn't masked by a cleanup error on the
return path.
+ private void releaseSegmentRecordBuffer() {
+ if (_segmentRecordBuffer != null) {
+ try {
+ _segmentRecordBuffer.close();
+ } catch (IOException e) {
+ LOGGER.warn("Failed to close segment record buffer", e);
}
- };
+ _segmentRecordBuffer = null;
+ }
+ if (_segmentRecordFile.exists()) {
+ try {
+ FileUtils.forceDelete(_segmentRecordFile);
+ } catch (IOException e) {
+ LOGGER.warn("Failed to delete segment record file: {}",
_segmentRecordFile, e);
+ }
+ }
}
@Override
@@ -352,7 +382,11 @@ public class OffHeapSingleTreeBuilder extends
BaseSingleTreeBuilder {
@Override
public void close()
throws IOException {
- super.close();
+ try {
+ super.close();
+ } finally {
+ releaseSegmentRecordBuffer();
+ }
if (_starTreeRecordBuffer != null) {
_starTreeRecordBuffer.close();
}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java
index 19e2cba5e31..33da25b4520 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/startree/v2/builder/OffHeapSingleTreeBuilderTest.java
@@ -18,16 +18,63 @@
*/
package org.apache.pinot.segment.local.startree.v2.builder;
+import java.io.File;
+import java.util.Arrays;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.TreeMap;
+import org.apache.commons.configuration2.PropertiesConfiguration;
+import org.apache.commons.io.FileUtils;
+import
org.apache.pinot.segment.local.indexsegment.immutable.ImmutableSegmentLoader;
+import
org.apache.pinot.segment.local.segment.creator.impl.SegmentIndexCreationDriverImpl;
+import org.apache.pinot.segment.local.segment.readers.GenericRowRecordReader;
+import org.apache.pinot.segment.local.startree.StarTreeBuilderUtils;
import
org.apache.pinot.segment.local.startree.v2.builder.OffHeapSingleTreeBuilder.FixedSizeRecordOffsets;
import
org.apache.pinot.segment.local.startree.v2.builder.OffHeapSingleTreeBuilder.RecordOffsets;
import
org.apache.pinot.segment.local.startree.v2.builder.OffHeapSingleTreeBuilder.VariableSizeRecordOffsets;
+import org.apache.pinot.segment.spi.ImmutableSegment;
+import org.apache.pinot.segment.spi.creator.SegmentGeneratorConfig;
+import org.apache.pinot.segment.spi.index.reader.Dictionary;
+import org.apache.pinot.segment.spi.index.reader.ForwardIndexReader;
+import org.apache.pinot.segment.spi.index.reader.ForwardIndexReaderContext;
+import org.apache.pinot.segment.spi.index.startree.StarTreeV2;
+import org.apache.pinot.segment.spi.memory.PinotDataBuffer;
+import org.apache.pinot.spi.config.table.StarTreeIndexConfig;
+import org.apache.pinot.spi.config.table.TableConfig;
+import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.data.readers.GenericRow;
+import org.apache.pinot.spi.utils.ReadMode;
+import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
+import org.testng.annotations.AfterMethod;
+import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertTrue;
public class OffHeapSingleTreeBuilderTest {
+ private static final File TEMP_DIR = new File(FileUtils.getTempDirectory(),
"OffHeapSingleTreeBuilderTest");
+ private static final File INDEX_DIR = new File(TEMP_DIR, "testSegment");
+ private static final String SEGMENT_RECORD_BUFFER_DESCRIPTION =
"OffHeapSingleTreeBuilder: segment record buffer";
+
+ @BeforeMethod
+ public void setUp()
+ throws Exception {
+ FileUtils.deleteQuietly(TEMP_DIR);
+ FileUtils.forceMkdir(TEMP_DIR);
+ }
+
+ @AfterMethod
+ public void tearDown() {
+ FileUtils.deleteQuietly(TEMP_DIR);
+ }
+
@Test
public void testFixedSizeRecordOffsets() {
RecordOffsets offsets = new FixedSizeRecordOffsets(1 << 30);
@@ -53,4 +100,211 @@ public class OffHeapSingleTreeBuilderTest {
assertEquals(offsets.getStartOffset(3), Integer.MAX_VALUE + 456L);
assertEquals(offsets.getEndOffset(), Integer.MAX_VALUE + 456L + 789L);
}
+
+ /// Drives `sortAndAggregateSegmentRecords` on a real segment and asserts
the segment record
+ /// buffer allocated in Sub-phase A is released once the iterator is
drained. Targets our
+ /// specific buffer by description via `PinotDataBuffer.getBufferInfo()` so
the assertion is
+ /// invariant under any concurrent JVM-wide buffer traffic.
+ @Test
+ public void testSegmentRecordBufferReleasedAfterIteratorDrain()
+ throws Exception {
+ buildTestSegment();
+ File segmentDir = INDEX_DIR.listFiles()[0];
+ ImmutableSegment segment = ImmutableSegmentLoader.load(segmentDir,
ReadMode.mmap);
+ try {
+ List<StarTreeV2BuilderConfig> builderConfigs =
createBuilderConfigs(segment);
+ File outputDir = new File(TEMP_DIR, "starTreeOutputDrain");
+ FileUtils.forceMkdir(outputDir);
+ assertFalse(isSegmentRecordBufferLive(), "No pre-existing segment record
buffer");
+
+ try (OffHeapSingleTreeBuilder builder = new
OffHeapSingleTreeBuilder(builderConfigs.get(0), outputDir, segment,
+ new PropertiesConfiguration())) {
+ int numDocs = segment.getSegmentMetadata().getTotalDocs();
+ Iterator<?> iterator = builder.sortAndAggregateSegmentRecords(numDocs);
+ assertTrue(isSegmentRecordBufferLive(),
+ "Segment record buffer should be live after
sortAndAggregateSegmentRecords");
+
+ while (iterator.hasNext()) {
+ iterator.next();
+ }
+
+ assertFalse(isSegmentRecordBufferLive(), "Segment record buffer should
be released after iterator drain");
+ }
+ assertFalse(isSegmentRecordBufferLive(), "Segment record buffer should
remain released after close()");
+ } finally {
+ segment.destroy();
+ }
+ }
+
+ /// Closes the builder without draining the iterator —
`releaseSegmentRecordBuffer()` on the
+ /// close path must still release the buffer. Prior to the fix, `close()`
did not release
+ /// `_segmentRecordBuffer`, so an operator abandoning a partial build would
leak direct memory
+ /// until the JVM cleaner ran (unbounded).
+ @Test
+ public void testSegmentRecordBufferReleasedOnCloseWithoutDrain()
+ throws Exception {
+ buildTestSegment();
+ File segmentDir = INDEX_DIR.listFiles()[0];
+ ImmutableSegment segment = ImmutableSegmentLoader.load(segmentDir,
ReadMode.mmap);
+ try {
+ List<StarTreeV2BuilderConfig> builderConfigs =
createBuilderConfigs(segment);
+ File outputDir = new File(TEMP_DIR, "starTreeOutputAbandon");
+ FileUtils.forceMkdir(outputDir);
+ assertFalse(isSegmentRecordBufferLive(), "No pre-existing segment record
buffer");
+
+ OffHeapSingleTreeBuilder builder = new
OffHeapSingleTreeBuilder(builderConfigs.get(0), outputDir, segment,
+ new PropertiesConfiguration());
+ try {
+ int numDocs = segment.getSegmentMetadata().getTotalDocs();
+ Iterator<?> iterator = builder.sortAndAggregateSegmentRecords(numDocs);
+ // Pull exactly one element; leave the iterator undrained so close()
must release the buffer.
+ iterator.next();
+ assertTrue(isSegmentRecordBufferLive(), "Segment record buffer should
be live mid-iteration");
+ } finally {
+ builder.close();
+ }
+
+ assertFalse(isSegmentRecordBufferLive(),
+ "Segment record buffer should be released after close() on undrained
iterator");
+ } finally {
+ segment.destroy();
+ }
+ }
+
+ private static boolean isSegmentRecordBufferLive() {
+ return PinotDataBuffer.getBufferInfo().stream().anyMatch(info ->
info.contains(SEGMENT_RECORD_BUFFER_DESCRIPTION));
+ }
+
+ /// 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. Direct regression test for 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];
+
+ File offHeapDir = new File(TEMP_DIR, "offHeapCopy");
+ File onHeapDir = new File(TEMP_DIR, "onHeapCopy");
+ FileUtils.copyDirectory(sourceSegmentDir, offHeapDir);
+ FileUtils.copyDirectory(sourceSegmentDir, onHeapDir);
+
+ List<StarTreeV2BuilderConfig> offHeapBuilderConfigs;
+ List<StarTreeV2BuilderConfig> onHeapBuilderConfigs;
+ ImmutableSegment sourceSegment =
ImmutableSegmentLoader.load(sourceSegmentDir, ReadMode.mmap);
+ try {
+ offHeapBuilderConfigs = createBuilderConfigs(sourceSegment);
+ onHeapBuilderConfigs = createBuilderConfigs(sourceSegment);
+ } finally {
+ sourceSegment.destroy();
+ }
+
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(offHeapBuilderConfigs, offHeapDir,
+ MultipleTreesBuilder.BuildMode.OFF_HEAP)) {
+ builder.build();
+ }
+ try (MultipleTreesBuilder builder = new
MultipleTreesBuilder(onHeapBuilderConfigs, 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", DataType.STRING)
+ .addSingleValueDimension("intCol", DataType.INT)
+ .addMetric("longCol", DataType.LONG)
+ .build();
+
+ TableConfig tableConfig = new TableConfigBuilder(TableType.OFFLINE)
+ .setTableName("testTable")
+ .build();
+
+ SegmentGeneratorConfig config = new SegmentGeneratorConfig(tableConfig,
schema);
+ config.setOutDir(TEMP_DIR.getAbsolutePath());
+ config.setSegmentName("testSegment");
+
+ // Produce a few rows with distinct dim tuples so Sub-phase C's
sortedDocIds order genuinely
+ // differs from segment docId order — exercises the "reads dims from
buffer under sorted docId
+ // access" path.
+ List<GenericRow> rows = Arrays.asList(
+ createRow("A", 1, 10L),
+ createRow("B", 2, 20L),
+ createRow("A", 2, 30L),
+ createRow("B", 1, 40L),
+ createRow("C", 3, 50L)
+ );
+
+ SegmentIndexCreationDriverImpl driver = new
SegmentIndexCreationDriverImpl();
+ driver.init(config, new GenericRowRecordReader(rows));
+ driver.build();
+ }
+
+ private GenericRow createRow(String stringValue, int intValue, long
longValue) {
+ GenericRow row = new GenericRow();
+ row.putValue("stringCol", stringValue);
+ row.putValue("intCol", intValue);
+ row.putValue("longCol", longValue);
+ return row;
+ }
+
+ private List<StarTreeV2BuilderConfig> createBuilderConfigs(ImmutableSegment
segment)
+ throws Exception {
+ StarTreeIndexConfig starTreeConfig = new StarTreeIndexConfig(
+ Arrays.asList("stringCol", "intCol"),
+ null,
+ Arrays.asList("SUM__longCol"),
+ null,
+ 1000);
+ return StarTreeBuilderUtils.generateBuilderConfigs(
+ Arrays.asList(starTreeConfig),
+ false,
+ segment.getSegmentMetadata());
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]