This is an automated email from the ASF dual-hosted git repository.
xiangfu0 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 0f1c98d9e14 Add codec pipeline integration tests (#19309)
0f1c98d9e14 is described below
commit 0f1c98d9e147d1f2067add74af10227eff3c4372
Author: Xiang Fu <[email protected]>
AuthorDate: Wed Sep 30 02:58:14 2026 +0700
Add codec pipeline integration tests (#19309)
End-to-end coverage for every raw forward-index encoding Pinot can put on a
single-value column, on an OFFLINE and a REALTIME table: that the
configuration
survives the controller, reaches segment generation, selects the on-disk
format
it asked for, and reads back identically through both query engines.
CodecPipelineIntegrationTest is driven entirely by one declarative matrix of
(column, data type, encoding variant) rows. The schema, the Avro schema and
records, noDictionaryColumns, the FieldConfig list, the queries and every
assertion are all generated by iterating that list, so the configuration
under
test and the assertions about it cannot drift apart. Column names are
derived
from (type, variant): intLegacyPassThrough, longSpecDeltaT64Lz4, and so on.
One codec combination is one column. Per data type the matrix starts from a
PASS_THROUGH baseline - raw with no compression, the only way to say
"uncompressed", since the codec DSL has no identity codec - and then adds
one
column per combination, 68 in total:
- INT, LONG, FLOAT, DOUBLE, STRING and BYTES each get all five
raw-applicable
legacy compressionCodec values: PASS_THROUGH, SNAPPY, ZSTANDARD, LZ4,
GZIP.
BIG_DECIMAL is left out because it is stored as BYTES and adds no
forward-index shape; CLP* is a whole-index STRING format rather than a
chunk
codec, and DELTA/DELTADELTA are rejected on any column with a forward
index.
- INT and LONG additionally get 19 codecSpec pipelines each, since
ForwardIndexType.validateCodecPipelineShape accepts codecSpec only for
single-value INT/LONG columns. Applying the same 19 to both types
covers the
separate int and long implementations inside DELTA, DELTADELTA, T64 and
GORILLA.
The 19 pipelines exercise all eight built-ins and every shape the validator
accepts: each compression alone, ZSTD with an explicit non-default level,
each
transform alone, chained value-preserving transforms (DELTA,DELTADELTA), a
packing transform after a value-preserving one (DELTA,T64 and
DELTADELTA,GORILLA), transform + compression pairs covering all four
compressions, three-stage chains (DELTA,T64,LZ4 and
DELTADELTA,GORILLA,ZSTD(3)),
and a multi-compression chain (LZ4,GZIP). GORILLA,ZSTANDARD also covers the
legacy-spelling alias, whose canonical form is GORILLA,ZSTD(3).
The central invariant is that codecs are transparent: for a given data type,
every column must read back exactly what the baseline column reads back, on
every row. That is asserted with one whole-table query per (data type,
engine) -
one predicate per column, not one query per column - so the matrix stays
wide
while the query count stays small. Two more queries per data type read the
same
columns through the other two paths: a projection spot check on the rows
straddling the configured targetDocsPerChunk plus the first and last rows,
and
MIN/MAX per column through the aggregation path. Values are a pure function
of
the doc id, so all columns of a type agree by construction; the sequences
are
slowly increasing and cross zero, which is what
DELTA/DELTADELTA/T64/GORILLA are
designed for, and four doc ids carry the type's extremes instead.
testGeneratedSegmentsUseConfiguredForwardIndexFormats reads the generated
segments and asserts, for every matrix column, that the on-disk format
matches
its configuration: a codecSpec column must be the self-describing V7 reader
with
the exact canonical spec in its header and no legacy ChunkCompressionType,
and a
legacy column must not be V7 and must report exactly the
ChunkCompressionType
its codec maps to - including PASS_THROUGH, so a regression that routed
every raw
column through V7, or compressed everything anyway, cannot pass.
CodecPipelineRealtimeIntegrationTest stays deliberately narrow and does not
repeat the matrix: once a segment is on disk, the codecs and the read paths
are
table-type agnostic. It covers only what is realtime-specific - the
configured
format is ignored while rows are consuming and applied when a consuming
segment
is committed - with one LONG codecSpec column, one INT codecSpec column and
a
legacy PASS_THROUGH column as the negative control for the V7 routing
branch.
---
.../tests/custom/CodecPipelineIntegrationTest.java | 973 +++++++++++++++++++++
.../CodecPipelineRealtimeIntegrationTest.java | 432 +++++++++
2 files changed, 1405 insertions(+)
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineIntegrationTest.java
new file mode 100644
index 00000000000..37aa6dee244
--- /dev/null
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineIntegrationTest.java
@@ -0,0 +1,973 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.integration.tests.custom;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.node.ObjectNode;
+import java.io.File;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import javax.annotation.Nullable;
+import org.apache.avro.Schema.Field;
+import org.apache.avro.Schema.Type;
+import org.apache.avro.file.DataFileWriter;
+import org.apache.avro.generic.GenericData;
+import
org.apache.pinot.segment.local.segment.index.readers.forward.FixedByteChunkSVForwardIndexReaderV7;
+import org.apache.pinot.segment.local.segment.store.SegmentLocalFSDirectory;
+import org.apache.pinot.segment.spi.ColumnMetadata;
+import org.apache.pinot.segment.spi.compression.ChunkCompressionType;
+import org.apache.pinot.segment.spi.index.FieldIndexConfigs;
+import org.apache.pinot.segment.spi.index.ForwardIndexConfig;
+import org.apache.pinot.segment.spi.index.StandardIndexes;
+import org.apache.pinot.segment.spi.index.reader.ForwardIndexReader;
+import org.apache.pinot.segment.spi.memory.PinotDataBuffer;
+import org.apache.pinot.segment.spi.store.SegmentDirectory;
+import org.apache.pinot.spi.config.table.FieldConfig;
+import org.apache.pinot.spi.config.table.FieldConfig.CompressionCodec;
+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.utils.BytesUtils;
+import org.apache.pinot.spi.utils.JsonUtils;
+import org.apache.pinot.spi.utils.ReadMode;
+import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertTrue;
+
+
+/// Integration test for raw forward-index encodings on an OFFLINE table,
covering every encoding and
+/// compression codec Pinot can put on a raw single-value column.
+///
+/// Everything is derived from one declarative matrix, [#MATRIX]: a list of
+/// (column, data type, encoding variant) rows where the variant is either a
legacy
+/// [CompressionCodec] or a `codecSpec` pipeline string. The schema, the Avro
schema, the Avro record
+/// population, `noDictionaryColumns`, the `FieldConfig` list, every query and
every assertion are
+/// generated by iterating that one list, so the configuration under test and
the assertions about it
+/// cannot drift apart. Column names are derived from (type, variant), e.g.
`intLegacyPassThrough`
+/// and `longSpecDeltaT64Lz4`.
+///
+/// **One codec combination is one column.** Per data type the matrix starts
from a
+/// `PASS_THROUGH` baseline — raw with no compression, the only way to say
"uncompressed", since the
+/// codec DSL has no identity codec — and then adds one column per codec
combination:
+///
+/// - Every data type gets the four remaining raw-applicable legacy codecs
(`SNAPPY`, `ZSTANDARD`,
+/// `LZ4`, `GZIP`). `CLP*` is a whole-index STRING format rather than a
chunk codec and
+/// `DELTA`/`DELTADELTA` are rejected on any column with a forward index, so
neither belongs here.
+/// - INT and LONG additionally get one column per `codecSpec` in
[#CODEC_SPECS]. They are the only
+/// types that may: `ForwardIndexType.validateCodecPipelineShape` accepts
`codecSpec` only for
+/// single-value INT/LONG columns, and `ForwardIndexConfig` rejects
`codecSpec` together with
+/// `compressionCodec`. Every other type can therefore only vary by legacy
codec.
+///
+/// The central invariant is that **codecs are transparent**: for a given data
type, every column
+/// must read back exactly what the baseline column reads back, on every row.
That is asserted with
+/// one whole-table query per (data type, engine) rather than one query per
column, so the matrix can
+/// stay wide while the query count stays small.
+///
+/// `FieldConfig.encodingType` + `indexes.forward.targetDocsPerChunk` is set
on every matrix column
+/// and `compressionCodec` is set at the `FieldConfig` level, which also
exercises the deserializer
+/// branch in `ForwardIndexType` that reconciles the legacy column-level
signal with the modern
+/// `indexes.forward` block.
+///
+/// Codec arithmetic, corrupt input, and the V7 writer/reader round trip are
unit tested in
+/// `pinot-segment-local`; what only an integration test can cover is that the
configuration survives
+/// the controller, reaches segment generation, selects the format it asked
for, and reads back
+/// identically through both query engines.
+@Test(suiteName = "CustomClusterIntegrationTest")
+public class CodecPipelineIntegrationTest extends
CustomDataQueryClusterIntegrationTest {
+
+ private static final String TABLE_NAME = "CodecPipelineIntegrationTest";
+ private static final int NUM_DOCS = 1000;
+ /// Each of the two Avro files holds 500 rows, so 256 puts a real chunk
boundary inside every
+ /// generated segment for the point lookups below to straddle.
+ private static final int TARGET_DOCS_PER_CHUNK = 256;
+
+ private static final String TIME_COL = "ts";
+ /// Dictionary-encoded column: proves dictionary and raw codec columns
coexist in one segment.
+ private static final String DICT_STR_COL = "dictStr";
+
+ /// Doc ids that carry a deliberately hostile value instead of the slowly
increasing sequence: the
+ /// type's smallest and largest value, zero, and minus one for the numeric
types, and the shortest,
+ /// longest, all-zero and non-ASCII value for STRING and BYTES. The value is
still a pure function of
+ /// the doc id, so every column of a type still agrees on these rows.
+ private static final int MIN_VALUE_DOC_ID = 100;
+ private static final int MAX_VALUE_DOC_ID = 200;
+ private static final int ZERO_VALUE_DOC_ID = 300;
+ private static final int NEG_ONE_VALUE_DOC_ID = 400;
+
+ /// Slowly increasing sequences: what DELTA, DELTADELTA, T64 and GORILLA are
designed for. Each one
+ /// starts negative and crosses zero, so negative values are covered outside
the hostile doc ids too.
+ private static final int INT_BASE = -3000;
+ private static final int INT_STEP = 7;
+ private static final long LONG_BASE = -2_000_000_000_000L;
+ private static final long LONG_STEP = 3_000_000_001L;
+ private static final float FLOAT_BASE = -1234.5f;
+ private static final float FLOAT_STEP = 2.5f;
+ private static final double DOUBLE_BASE = -1234.25d;
+ private static final double DOUBLE_STEP = 2.5d;
+
+ /// Raw single-value types this fixture drives end to end. BIG_DECIMAL is
deliberately absent: it
+ /// is stored as BYTES, so it adds no forward-index shape, and Avro has no
decimal primitive.
+ private static final List<DataType> DATA_TYPES =
+ List.of(DataType.INT, DataType.LONG, DataType.FLOAT, DataType.DOUBLE,
DataType.STRING, DataType.BYTES);
+
+ /// The uncompressed baseline every other column of the same type is
compared against.
+ private static final CompressionCodec BASELINE_CODEC =
CompressionCodec.PASS_THROUGH;
+
+ /// Every legacy [CompressionCodec] with `isApplicableToRawIndex() == true`,
baseline first.
+ private static final List<CompressionCodec> LEGACY_CODECS =
+ List.of(BASELINE_CODEC, CompressionCodec.SNAPPY,
CompressionCodec.ZSTANDARD, CompressionCodec.LZ4,
+ CompressionCodec.GZIP);
+
+ /// One entry per `codecSpec` combination, applied to both INT and LONG so
the separate int and long
+ /// implementations inside DELTA, DELTADELTA, T64 and GORILLA are both
covered.
+ ///
+ /// Collectively these exercise all eight built-ins (DELTA, DELTADELTA, T64,
GORILLA, ZSTD, LZ4,
+ /// SNAPPY, GZIP) and every shape the pipeline validator accepts:
compression alone, a transform
+ /// alone, chained value-preserving transforms, a packing transform after a
value-preserving one,
+ /// transform + compression (covering all four compressions between them),
three-stage chains, and a
+ /// multi-compression chain.
+ private static final List<CodecSpecVariant> CODEC_SPECS = List.of(
+ // Compression only.
+ new CodecSpecVariant("LZ4", "LZ4"),
+ new CodecSpecVariant("SNAPPY", "SNAPPY"),
+ new CodecSpecVariant("GZIP", "GZIP"),
+ // ZSTD's default level is materialized into the canonical spec written
to the segment header.
+ new CodecSpecVariant("ZSTD", "ZSTD(3)"),
+ new CodecSpecVariant("ZSTD(9)", "ZSTD(9)"),
+ // Transform only. DELTA/DELTADELTA preserve the typed value layout;
T64/GORILLA pack it.
+ new CodecSpecVariant("DELTA", "DELTA"),
+ new CodecSpecVariant("DELTADELTA", "DELTADELTA"),
+ new CodecSpecVariant("T64", "T64"),
+ new CodecSpecVariant("GORILLA", "GORILLA"),
+ // Chained value-preserving transforms, and a packing transform after a
value-preserving one.
+ new CodecSpecVariant("DELTA,DELTADELTA", "DELTA,DELTADELTA"),
+ new CodecSpecVariant("DELTA,T64", "DELTA,T64"),
+ new CodecSpecVariant("DELTADELTA,GORILLA", "DELTADELTA,GORILLA"),
+ // Transform + compression, covering all four compressions across the
four rows.
+ new CodecSpecVariant("DELTA,LZ4", "DELTA,LZ4"),
+ new CodecSpecVariant("DELTADELTA,SNAPPY", "DELTADELTA,SNAPPY"),
+ new CodecSpecVariant("T64,GZIP", "T64,GZIP"),
+ // ZSTANDARD is accepted as a legacy-spelling alias; canonicalization
still emits ZSTD.
+ new CodecSpecVariant("GORILLA,ZSTANDARD", "GORILLA,ZSTD(3)"),
+ // Three stages: value-preserving transform, packing transform,
compression.
+ new CodecSpecVariant("DELTA,T64,LZ4", "DELTA,T64,LZ4"),
+ new CodecSpecVariant("DELTADELTA,GORILLA,ZSTD(3)",
"DELTADELTA,GORILLA,ZSTD(3)"),
+ // Several compression stages chained, with no transform at all.
+ new CodecSpecVariant("LZ4,GZIP", "LZ4,GZIP"));
+
+ /// Data types that may carry a `codecSpec`, per
`ForwardIndexType.validateCodecPipelineShape`.
+ private static final List<DataType> CODEC_SPEC_TYPES = List.of(DataType.INT,
DataType.LONG);
+
+ /// The single source of truth: one row per column under test.
+ private static final List<MatrixColumn> MATRIX = buildMatrix();
+
+ /// First and last row, the rows straddling the [#TARGET_DOCS_PER_CHUNK]
boundary in both generated
+ /// segments (local rows 255/256 are doc ids 510/512 in one segment and
511/513 in the other), and
+ /// the four hostile-value rows.
+ private static final int[] SPOT_CHECK_DOC_IDS =
+ ascending(0, 1, MIN_VALUE_DOC_ID, MAX_VALUE_DOC_ID, ZERO_VALUE_DOC_ID,
NEG_ONE_VALUE_DOC_ID, 510, 511, 512, 513,
+ 998, 999);
+
+ /// The spot-check doc ids as a SQL `IN` list. Kept sorted so it lines up
with `ORDER BY ts`.
+ private static final String SPOT_CHECK_DOC_ID_LIST =
docIdList(SPOT_CHECK_DOC_IDS);
+
+ /// Sorted so the rows come back in the same order the assertions walk them.
+ private static int[] ascending(int... docIds) {
+ int[] sorted = docIds.clone();
+ Arrays.sort(sorted);
+ return sorted;
+ }
+
+ private static String docIdList(int[] docIds) {
+ StringBuilder builder = new StringBuilder();
+ for (int docId : docIds) {
+ if (builder.length() > 0) {
+ builder.append(", ");
+ }
+ builder.append(docId);
+ }
+ return builder.toString();
+ }
+
+ /// A `codecSpec` combination: the string written into the table config, and
the canonical form the
+ /// codec runtime derives from it and freezes into the V7 segment header.
+ private static final class CodecSpecVariant {
+ final String _configured;
+ final String _canonical;
+
+ private CodecSpecVariant(String configured, String canonical) {
+ _configured = configured;
+ _canonical = canonical;
+ }
+ }
+
+ /// One column of the matrix. Exactly one of [#_compressionCodec] and
[#_codecSpec] is set, because
+ /// `ForwardIndexConfig` rejects both at once.
+ private static final class MatrixColumn {
+ final String _column;
+ final DataType _dataType;
+ @Nullable
+ final CompressionCodec _compressionCodec;
+ @Nullable
+ final CodecSpecVariant _codecSpec;
+
+ private MatrixColumn(String column, DataType dataType, @Nullable
CompressionCodec compressionCodec,
+ @Nullable CodecSpecVariant codecSpec) {
+ _column = column;
+ _dataType = dataType;
+ _compressionCodec = compressionCodec;
+ _codecSpec = codecSpec;
+ }
+
+ boolean isCodecPipeline() {
+ return _codecSpec != null;
+ }
+
+ /// The configured variant, for assertion messages.
+ String variant() {
+ return _codecSpec != null ? "codecSpec=" + _codecSpec._configured :
"compressionCodec=" + _compressionCodec;
+ }
+
+ @Override
+ public String toString() {
+ return _column + "[" + _dataType + ", " + variant() + "]";
+ }
+ }
+
+ /// Builds the matrix: per data type the legacy codecs, then the `codecSpec`
combinations for the
+ /// types that accept them. Generated names must be unique, which a
duplicate key would break
+ /// silently (two matrix rows pointing at one physical column), so that is
checked here.
+ private static List<MatrixColumn> buildMatrix() {
+ Map<String, MatrixColumn> byName = new LinkedHashMap<>();
+ for (DataType dataType : DATA_TYPES) {
+ for (CompressionCodec codec : LEGACY_CODECS) {
+ String name = typePrefix(dataType) + "Legacy" +
camelSuffix(codec.name());
+ MatrixColumn previous = byName.put(name, new MatrixColumn(name,
dataType, codec, null));
+ if (previous != null) {
+ throw new IllegalStateException("Duplicate generated column name: "
+ name);
+ }
+ }
+ if (!CODEC_SPEC_TYPES.contains(dataType)) {
+ continue;
+ }
+ for (CodecSpecVariant variant : CODEC_SPECS) {
+ String name = typePrefix(dataType) + "Spec" +
camelSuffix(variant._configured);
+ MatrixColumn previous = byName.put(name, new MatrixColumn(name,
dataType, null, variant));
+ if (previous != null) {
+ throw new IllegalStateException("Duplicate generated column name: "
+ name);
+ }
+ }
+ }
+ return List.copyOf(byName.values());
+ }
+
+ private static String typePrefix(DataType dataType) {
+ return dataType.name().toLowerCase(Locale.ROOT);
+ }
+
+ /// `PASS_THROUGH` to `PassThrough`, `DELTA,T64,LZ4` to `DeltaT64Lz4`,
`ZSTD(9)` to `Zstd9`: every
+ /// run of non-alphanumeric characters starts a new capitalized token, so
the name is a
+ /// deterministic function of the variant and stays a legal Pinot column
name.
+ private static String camelSuffix(String raw) {
+ StringBuilder builder = new StringBuilder(raw.length());
+ boolean startOfToken = true;
+ for (int i = 0; i < raw.length(); i++) {
+ char c = raw.charAt(i);
+ if (Character.isLetterOrDigit(c)) {
+ builder.append(startOfToken ? Character.toUpperCase(c) :
Character.toLowerCase(c));
+ startOfToken = false;
+ } else {
+ startOfToken = true;
+ }
+ }
+ return builder.toString();
+ }
+
+ private static List<MatrixColumn> columnsOfType(DataType dataType) {
+ List<MatrixColumn> columns = new ArrayList<>();
+ for (MatrixColumn column : MATRIX) {
+ if (column._dataType == dataType) {
+ columns.add(column);
+ }
+ }
+ return columns;
+ }
+
+ private static MatrixColumn baselineOfType(DataType dataType) {
+ for (MatrixColumn column : MATRIX) {
+ if (column._dataType == dataType && column._compressionCodec ==
BASELINE_CODEC) {
+ return column;
+ }
+ }
+ throw new IllegalStateException("No " + BASELINE_CODEC + " baseline column
for " + dataType);
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Values. Every value is a pure function of the doc id, so for a given row
all columns of a data
+ // type hold the same value and the cross-codec equality assertions below
are exact.
+ //
---------------------------------------------------------------------------------------------
+
+ private static int intValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return Integer.MIN_VALUE;
+ case MAX_VALUE_DOC_ID:
+ return Integer.MAX_VALUE;
+ case ZERO_VALUE_DOC_ID:
+ return 0;
+ case NEG_ONE_VALUE_DOC_ID:
+ return -1;
+ default:
+ return INT_BASE + docId * INT_STEP;
+ }
+ }
+
+ private static long longValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return Long.MIN_VALUE;
+ case MAX_VALUE_DOC_ID:
+ return Long.MAX_VALUE;
+ case ZERO_VALUE_DOC_ID:
+ return 0L;
+ case NEG_ONE_VALUE_DOC_ID:
+ return -1L;
+ default:
+ return LONG_BASE + docId * LONG_STEP;
+ }
+ }
+
+ private static float floatValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return -Float.MAX_VALUE;
+ case MAX_VALUE_DOC_ID:
+ return Float.MAX_VALUE;
+ case ZERO_VALUE_DOC_ID:
+ return 0.0f;
+ case NEG_ONE_VALUE_DOC_ID:
+ return -1.0f;
+ default:
+ // Every value is a small multiple of 0.5, which is exact in float, so
the sequence survives the
+ // Avro and JSON round trips bit for bit and can be compared with
floatToIntBits.
+ return FLOAT_BASE + docId * FLOAT_STEP;
+ }
+ }
+
+ private static double doubleValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return -Double.MAX_VALUE;
+ case MAX_VALUE_DOC_ID:
+ return Double.MAX_VALUE;
+ case ZERO_VALUE_DOC_ID:
+ return 0.0d;
+ case NEG_ONE_VALUE_DOC_ID:
+ return -1.0d;
+ default:
+ return DOUBLE_BASE + docId * DOUBLE_STEP;
+ }
+ }
+
+ private static String stringValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return "-";
+ case MAX_VALUE_DOC_ID:
+ return "z".repeat(64);
+ case ZERO_VALUE_DOC_ID:
+ return "0";
+ case NEG_ONE_VALUE_DOC_ID:
+ // Multi-byte UTF-8, so the var-byte chunk length is not the character
count.
+ return "é中文";
+ default:
+ return "str-" + docId + "-" + "x".repeat(docId % 17);
+ }
+ }
+
+ private static byte[] bytesValueFor(int docId) {
+ switch (docId) {
+ case MIN_VALUE_DOC_ID:
+ return new byte[]{0};
+ case MAX_VALUE_DOC_ID:
+ byte[] longest = new byte[64];
+ Arrays.fill(longest, (byte) 0xFF);
+ return longest;
+ case ZERO_VALUE_DOC_ID:
+ return new byte[]{0, 0, 0, 0};
+ case NEG_ONE_VALUE_DOC_ID:
+ return new byte[]{(byte) 0xDE, (byte) 0xAD, (byte) 0xBE, (byte) 0xEF};
+ default:
+ byte[] value = new byte[1 + docId % 13];
+ for (int i = 0; i < value.length; i++) {
+ value[i] = (byte) (docId + i);
+ }
+ return value;
+ }
+ }
+
+ /// The Avro value for one column in one row. A fresh object per call, so no
two fields of a record
+ /// can share a mutable buffer.
+ private static Object avroValueFor(DataType dataType, int docId) {
+ switch (dataType) {
+ case INT:
+ return intValueFor(docId);
+ case LONG:
+ return longValueFor(docId);
+ case FLOAT:
+ return floatValueFor(docId);
+ case DOUBLE:
+ return doubleValueFor(docId);
+ case STRING:
+ return stringValueFor(docId);
+ case BYTES:
+ return ByteBuffer.wrap(bytesValueFor(docId));
+ default:
+ throw new IllegalStateException("No Avro value generator for " +
dataType);
+ }
+ }
+
+ private static Type avroTypeFor(DataType dataType) {
+ switch (dataType) {
+ case INT:
+ return Type.INT;
+ case LONG:
+ return Type.LONG;
+ case FLOAT:
+ return Type.FLOAT;
+ case DOUBLE:
+ return Type.DOUBLE;
+ case STRING:
+ return Type.STRING;
+ case BYTES:
+ return Type.BYTES;
+ default:
+ throw new IllegalStateException("No Avro type for " + dataType);
+ }
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Table, schema and data, all generated from MATRIX.
+ //
---------------------------------------------------------------------------------------------
+
+ @Override
+ public String getTableName() {
+ return TABLE_NAME;
+ }
+
+ @Override
+ protected long getCountStarResult() {
+ return NUM_DOCS;
+ }
+
+ @Override
+ public Schema createSchema() {
+ Schema.SchemaBuilder builder = new
Schema.SchemaBuilder().setSchemaName(getTableName());
+ for (MatrixColumn column : MATRIX) {
+ builder.addSingleValueDimension(column._column, column._dataType);
+ }
+ builder.addSingleValueDimension(DICT_STR_COL, DataType.STRING);
+ builder.addDateTimeField(TIME_COL, DataType.LONG, "1:MILLISECONDS:EPOCH",
"1:MILLISECONDS");
+ return builder.build();
+ }
+
+ @Override
+ public List<File> createAvroFiles()
+ throws IOException {
+ org.apache.avro.Schema avroSchema =
org.apache.avro.Schema.createRecord("codecRecord", null, null, false);
+ List<Field> fields = new ArrayList<>();
+ for (MatrixColumn column : MATRIX) {
+ fields.add(new Field(column._column,
org.apache.avro.Schema.create(avroTypeFor(column._dataType)), null, null));
+ }
+ fields.add(new Field(DICT_STR_COL,
org.apache.avro.Schema.create(Type.STRING), null, null));
+ fields.add(new Field(TIME_COL, org.apache.avro.Schema.create(Type.LONG),
null, null));
+ avroSchema.setFields(fields);
+
+ try (AvroFilesAndWriters avroFilesAndWriters =
createAvroFilesAndWriters(avroSchema)) {
+ List<DataFileWriter<GenericData.Record>> writers =
avroFilesAndWriters.getWriters();
+ for (int docId = 0; docId < NUM_DOCS; docId++) {
+ GenericData.Record record = new GenericData.Record(avroSchema);
+ for (MatrixColumn column : MATRIX) {
+ record.put(column._column, avroValueFor(column._dataType, docId));
+ }
+ record.put(DICT_STR_COL, "dict-" + docId);
+ record.put(TIME_COL, (long) docId);
+ writers.get(docId % getNumAvroFiles()).append(record);
+ }
+ return avroFilesAndWriters.getAvroFiles();
+ }
+ }
+
+ /// Built explicitly rather than through the base class helper: no sorted
column and no inverted,
+ /// range, or bloom index, so the forward index is the only thing under test
here.
+ @Override
+ public TableConfig createOfflineTableConfig() {
+ // No time column is set: `ts` carries the raw doc id, which is outside
Pinot's valid time
+ // interval, and segment generation would reject it if it were declared as
the time column.
+ return new
TableConfigBuilder(TableType.OFFLINE).setTableName(getTableName())
+ .setNoDictionaryColumns(getNoDictionaryColumns())
+ .setFieldConfigList(getFieldConfigs())
+ .build();
+ }
+
+ @Override
+ protected List<String> getNoDictionaryColumns() {
+ // DICT_STR_COL keeps its dictionary, so it is intentionally NOT in this
list.
+ List<String> noDictionaryColumns = new ArrayList<>(MATRIX.size());
+ for (MatrixColumn column : MATRIX) {
+ noDictionaryColumns.add(column._column);
+ }
+ return noDictionaryColumns;
+ }
+
+ @Override
+ protected List<FieldConfig> getFieldConfigs() {
+ List<FieldConfig> fieldConfigs = new ArrayList<>(MATRIX.size() + 1);
+ for (MatrixColumn column : MATRIX) {
+ fieldConfigs.add(rawFieldConfig(column));
+ }
+ fieldConfigs.add(new FieldConfig.Builder(DICT_STR_COL)
+ .withEncodingType(FieldConfig.EncodingType.DICTIONARY)
+ .build());
+ return fieldConfigs;
+ }
+
+ /// A RAW `FieldConfig` for one matrix column. `codecSpec` only exists in
the modern
+ /// `indexes.forward` block (there is no top-level `FieldConfig.codecSpec`),
while the legacy codec
+ /// is set at the `FieldConfig` level; both carry `targetDocsPerChunk` in
`indexes.forward` so every
+ /// column has chunk boundaries inside every segment.
+ private static FieldConfig rawFieldConfig(MatrixColumn column) {
+ ObjectNode forward = JsonUtils.newObjectNode();
+ forward.put("targetDocsPerChunk", TARGET_DOCS_PER_CHUNK);
+ if (column.isCodecPipeline()) {
+ forward.put("codecSpec", column._codecSpec._configured);
+ }
+ ObjectNode indexes = JsonUtils.newObjectNode();
+ indexes.set("forward", forward);
+ FieldConfig.Builder builder = new FieldConfig.Builder(column._column)
+ .withEncodingType(FieldConfig.EncodingType.RAW)
+ .withIndexes(indexes);
+ if (!column.isCodecPipeline()) {
+ builder.withCompressionCodec(column._compressionCodec);
+ }
+ return builder.build();
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Query-level assertions.
+ //
---------------------------------------------------------------------------------------------
+
+ /// Cartesian product of (data type, query engine).
+ @DataProvider(name = "dataTypeAndEngine")
+ public Object[][] dataTypeAndEngine() {
+ List<Object[]> rows = new ArrayList<>(DATA_TYPES.size() * 2);
+ for (DataType dataType : DATA_TYPES) {
+ rows.add(new Object[]{dataType, false});
+ rows.add(new Object[]{dataType, true});
+ }
+ return rows.toArray(new Object[0][]);
+ }
+
+ /// The whole point of the matrix, for one data type on one engine, in three
queries:
+ ///
+ /// 1. Every column of the type equals the `PASS_THROUGH` baseline on
**every** row. One query with
+ /// one predicate per column, so a per-row decode error that preserves
aggregates cannot hide,
+ /// without paying a query per column. This reads through the filter/scan
path.
+ /// 1. The rows around a V7 chunk boundary, the first and last rows, and the
hostile-value rows are
+ /// fetched with every column of the type in the select list and compared
value by value against
+ /// the expected value. This reads through the projection path.
+ /// 1. For numeric types, `MIN`/`MAX` of every column — baseline included —
must equal the extremes
+ /// this test wrote. This reads through the aggregation path, and
asserting against the written
+ /// values rather than against the baseline also catches "every codec is
equally wrong". Pinot has
+ /// no `MIN`/`MAX` for STRING or BYTES, so that step is numeric only.
+ @Test(dataProvider = "dataTypeAndEngine")
+ public void testCodecMatrixPerDataType(DataType dataType, boolean
useMultiStageQueryEngine)
+ throws Exception {
+ setUseMultiStageQueryEngine(useMultiStageQueryEngine);
+ List<MatrixColumn> columns = columnsOfType(dataType);
+ MatrixColumn baseline = baselineOfType(dataType);
+ assertTrue(columns.size() > 1, "Expected more than one " + dataType + "
column in the matrix");
+
+ assertEveryRowAgreesWithBaseline(dataType, columns, baseline);
+ assertSpotCheckRows(dataType, columns);
+ if (dataType.isNumeric()) {
+ assertNumericAggregates(dataType, columns);
+ }
+ }
+
+ /// One predicate per non-baseline column: `<column> = <baseline>`, so a
single `COUNT(*)` states
+ /// that every column of the type agrees with the baseline on every row.
+ private static String agreementQuery(List<MatrixColumn> columns,
MatrixColumn baseline) {
+ String baselineExpression = comparisonExpression(baseline);
+ List<String> predicates = new ArrayList<>(columns.size() - 1);
+ for (MatrixColumn column : columns) {
+ if (column != baseline) {
+ predicates.add(comparisonExpression(column) + " = " +
baselineExpression);
+ }
+ }
+ return "SELECT COUNT(*) FROM " + TABLE_NAME + " WHERE " + String.join("
AND ", predicates);
+ }
+
+ /// Every column of the type plus `ts`, for the spot-check rows only,
ordered so the rows line up
+ /// with [#SPOT_CHECK_DOC_IDS].
+ private static String spotCheckQuery(List<MatrixColumn> columns) {
+ List<String> selectList = new ArrayList<>(columns.size() + 1);
+ selectList.add(TIME_COL);
+ for (MatrixColumn column : columns) {
+ selectList.add(column._column);
+ }
+ // An explicit LIMIT is required: without it the single-stage engine
applies its default
+ // selection limit of 10 and would silently return fewer rows than
SPOT_CHECK_DOC_IDS has.
+ return "SELECT " + String.join(", ", selectList) + " FROM " + TABLE_NAME +
" WHERE " + TIME_COL + " IN ("
+ + SPOT_CHECK_DOC_ID_LIST + ") ORDER BY " + TIME_COL + " LIMIT " +
SPOT_CHECK_DOC_IDS.length;
+ }
+
+ /// `MIN` and `MAX` for every column of the type, in one query, in column
order.
+ private static String minMaxQuery(List<MatrixColumn> columns) {
+ List<String> aggregates = new ArrayList<>(columns.size() * 2);
+ for (MatrixColumn column : columns) {
+ aggregates.add("MIN(" + column._column + ")");
+ aggregates.add("MAX(" + column._column + ")");
+ }
+ return "SELECT " + String.join(", ", aggregates) + " FROM " + TABLE_NAME;
+ }
+
+ private static String dictionarySpotCheckQuery() {
+ // See spotCheckQuery: the explicit LIMIT defeats the single-stage default
selection limit of 10.
+ return "SELECT " + TIME_COL + ", " + DICT_STR_COL + " FROM " + TABLE_NAME
+ " WHERE " + TIME_COL + " IN ("
+ + SPOT_CHECK_DOC_ID_LIST + ") ORDER BY " + TIME_COL + " LIMIT " +
SPOT_CHECK_DOC_IDS.length;
+ }
+
+ private static String dictionaryDistinctQuery() {
+ return "SELECT COUNT(DISTINCT " + DICT_STR_COL + ") FROM " + TABLE_NAME;
+ }
+
+ private void assertEveryRowAgreesWithBaseline(DataType dataType,
List<MatrixColumn> columns, MatrixColumn baseline)
+ throws Exception {
+ String query = agreementQuery(columns, baseline);
+ JsonNode result = postQuery(query);
+ assertEquals(result.get("resultTable").get("rows").get(0).get(0).asLong(),
NUM_DOCS,
+ "Not every " + dataType + " column agrees with the " + BASELINE_CODEC
+ " baseline on every row: " + query);
+ }
+
+ private void assertSpotCheckRows(DataType dataType, List<MatrixColumn>
columns)
+ throws Exception {
+ JsonNode result = postQuery(spotCheckQuery(columns));
+ JsonNode rows = result.get("resultTable").get("rows");
+ assertEquals(rows.size(), SPOT_CHECK_DOC_IDS.length, "Unexpected
spot-check row count for " + dataType);
+
+ for (int rowId = 0; rowId < SPOT_CHECK_DOC_IDS.length; rowId++) {
+ int docId = SPOT_CHECK_DOC_IDS[rowId];
+ JsonNode row = rows.get(rowId);
+ assertEquals(row.get(0).asInt(), docId, "Unexpected spot-check row order
for " + dataType);
+ for (int i = 0; i < columns.size(); i++) {
+ assertValueEquals(row.get(i + 1), dataType, docId, columns.get(i) + "
at ts=" + docId);
+ }
+ }
+ }
+
+ private void assertNumericAggregates(DataType dataType, List<MatrixColumn>
columns)
+ throws Exception {
+ JsonNode result = postQuery(minMaxQuery(columns));
+ JsonNode row = result.get("resultTable").get("rows").get(0);
+
+ // Every column, baseline included, is compared against the extremes this
test actually wrote, so
+ // the assertion catches both "one codec disagrees" and "every codec is
equally wrong".
+ if (dataType == DataType.LONG) {
+ // Compared as long, not double: this fixture writes Long.MIN_VALUE and
Long.MAX_VALUE, which
+ // are not exactly representable as double, so a double comparison would
tolerate an error of
+ // up to ~1024 at the extremes -- exactly where a bit-packing codec is
most likely to be wrong.
+ long expectedMin = longExtreme(false);
+ long expectedMax = longExtreme(true);
+ for (int i = 0; i < columns.size(); i++) {
+ MatrixColumn column = columns.get(i);
+ assertEquals(row.get(2 * i).asLong(), expectedMin, "Unexpected MIN(" +
column + ")");
+ assertEquals(row.get(2 * i + 1).asLong(), expectedMax, "Unexpected
MAX(" + column + ")");
+ }
+ return;
+ }
+ double expectedMin = extremeOf(dataType, false);
+ double expectedMax = extremeOf(dataType, true);
+ for (int i = 0; i < columns.size(); i++) {
+ MatrixColumn column = columns.get(i);
+ assertNumericEquals(row.get(2 * i), dataType, expectedMin, "Unexpected
MIN(" + column + ")");
+ assertNumericEquals(row.get(2 * i + 1), dataType, expectedMax,
"Unexpected MAX(" + column + ")");
+ }
+ }
+
+ /// Compares an aggregate result to an expected numeric value. FLOAT is
narrowed back to float
+ /// first: the two engines disagree on the aggregate's result type (one
widens to DOUBLE, the other
+ /// keeps FLOAT), and a float serialized as its own shortest decimal form
does not parse back to the
+ /// same double as the widened value. Narrowing to float is exact from
either representation.
+ private static void assertNumericEquals(JsonNode node, DataType dataType,
double expected, String context) {
+ if (dataType == DataType.FLOAT) {
+ assertEquals(Float.floatToIntBits((float) node.asDouble()),
Float.floatToIntBits((float) expected), context);
+ } else {
+ assertEquals(node.asDouble(), expected, context);
+ }
+ }
+
+ /// BYTES columns are compared through `bytesToHex` rather than directly:
hex equality is equivalent
+ /// for the invariant under test and keeps the predicate to string
comparison, which both engines
+ /// support identically. Every other type is compared as itself.
+ private static String comparisonExpression(MatrixColumn column) {
+ return column._dataType == DataType.BYTES ? "bytesToHex(" + column._column
+ ")" : column._column;
+ }
+
+ private static void assertValueEquals(JsonNode node, DataType dataType, int
docId, String context) {
+ switch (dataType) {
+ case INT:
+ assertEquals(node.asInt(), intValueFor(docId), context);
+ break;
+ case LONG:
+ assertEquals(node.asLong(), longValueFor(docId), context);
+ break;
+ case FLOAT:
+ // Bitwise comparison: the values written are exact in float, so
anything but an exact match
+ // is a decoding bug rather than a rounding artifact.
+ assertEquals(Float.floatToIntBits((float) node.asDouble()),
Float.floatToIntBits(floatValueFor(docId)),
+ context);
+ break;
+ case DOUBLE:
+ assertEquals(Double.doubleToLongBits(node.asDouble()),
Double.doubleToLongBits(doubleValueFor(docId)), context);
+ break;
+ case STRING:
+ assertEquals(node.asText(), stringValueFor(docId), context);
+ break;
+ case BYTES:
+ // Both engines render BYTES through BytesUtils.toHexString, so this
is an exact comparison.
+ assertEquals(node.asText(),
BytesUtils.toHexString(bytesValueFor(docId)), context);
+ break;
+ default:
+ throw new IllegalStateException("No value assertion for " + dataType);
+ }
+ }
+
+ /// The smallest or largest LONG this test writes, kept exact so the
extremes can be compared
+ /// without a lossy widening to double.
+ private static long longExtreme(boolean max) {
+ long result = max ? Long.MIN_VALUE : Long.MAX_VALUE;
+ for (int docId = 0; docId < NUM_DOCS; docId++) {
+ long value = longValueFor(docId);
+ result = max ? Math.max(result, value) : Math.min(result, value);
+ }
+ return result;
+ }
+
+ /// The smallest or largest value this test writes for a numeric type,
widened to double for
+ /// comparison against a `MIN`/`MAX` result.
+ private static double extremeOf(DataType dataType, boolean max) {
+ double result = max ? Double.NEGATIVE_INFINITY : Double.POSITIVE_INFINITY;
+ for (int docId = 0; docId < NUM_DOCS; docId++) {
+ double value;
+ switch (dataType) {
+ case INT:
+ value = intValueFor(docId);
+ break;
+ case LONG:
+ value = longValueFor(docId);
+ break;
+ case FLOAT:
+ value = floatValueFor(docId);
+ break;
+ case DOUBLE:
+ value = doubleValueFor(docId);
+ break;
+ default:
+ throw new IllegalStateException("Not a numeric type: " + dataType);
+ }
+ result = max ? Math.max(result, value) : Math.min(result, value);
+ }
+ return result;
+ }
+
+ /// A dictionary-encoded column in the same segment as every raw codec
column must still read back
+ /// correctly, which is what keeps the matrix from silently turning the
whole table raw.
+ @Test(dataProvider = "useBothQueryEngines")
+ public void testDictionaryColumnCoexists(boolean useMultiStageQueryEngine)
+ throws Exception {
+ setUseMultiStageQueryEngine(useMultiStageQueryEngine);
+
+ JsonNode result = postQuery(dictionarySpotCheckQuery());
+ JsonNode rows = result.get("resultTable").get("rows");
+ assertEquals(rows.size(), SPOT_CHECK_DOC_IDS.length, "Unexpected " +
DICT_STR_COL + " row count");
+ for (int rowId = 0; rowId < SPOT_CHECK_DOC_IDS.length; rowId++) {
+ assertEquals(rows.get(rowId).get(1).asText(), "dict-" +
SPOT_CHECK_DOC_IDS[rowId],
+ "Wrong " + DICT_STR_COL + " for ts=" + SPOT_CHECK_DOC_IDS[rowId]);
+ }
+
+ JsonNode distinct = postQuery(dictionaryDistinctQuery());
+
assertEquals(distinct.get("resultTable").get("rows").get(0).get(0).asLong(),
NUM_DOCS,
+ "Expected all " + NUM_DOCS + " distinct " + DICT_STR_COL + " values");
+ }
+
+ //
---------------------------------------------------------------------------------------------
+ // Format-level assertions: what actually landed on disk.
+ //
---------------------------------------------------------------------------------------------
+
+ /// Proves the configuration is not merely accepted and then ignored. Every
`codecSpec` column must
+ /// produce a self-describing V7 reader carrying the exact canonical spec
and no legacy
+ /// `ChunkCompressionType`; every legacy column must keep a non-V7 reader
and report exactly the
+ /// `ChunkCompressionType` its codec maps to — including `PASS_THROUGH`, so
an "everything is
+ /// compressed anyway" regression cannot pass. Derived from the same matrix,
so every column is
+ /// checked.
+ @Test
+ public void testGeneratedSegmentsUseConfiguredForwardIndexFormats()
+ throws Exception {
+ File[] segmentDirs = _segmentDir.listFiles(File::isDirectory);
+ assertNotNull(segmentDirs, "Segment output directory must be readable: " +
_segmentDir);
+ assertTrue(segmentDirs.length > 0, "Expected generated segments under " +
_segmentDir);
+
+ FieldIndexConfigs rawForwardConfig = new FieldIndexConfigs.Builder()
+ .add(StandardIndexes.forward(), new
ForwardIndexConfig.Builder(FieldConfig.EncodingType.RAW).build())
+ .build();
+ for (File segmentDir : segmentDirs) {
+ try (SegmentDirectory directory = new
SegmentLocalFSDirectory(segmentDir, ReadMode.mmap);
+ SegmentDirectory.Reader segmentReader = directory.createReader()) {
+ ColumnMetadata dictionaryMetadata =
directory.getSegmentMetadata().getColumnMetadataFor(DICT_STR_COL);
+ assertNotNull(dictionaryMetadata, "Missing metadata for " +
DICT_STR_COL + " in " + segmentDir);
+ assertTrue(dictionaryMetadata.hasDictionary(), DICT_STR_COL + " must
remain dictionary encoded");
+ assertEquals(dictionaryMetadata.getForwardIndexEncoding(),
FieldConfig.EncodingType.DICTIONARY);
+
+ for (MatrixColumn column : MATRIX) {
+ if (column.isCodecPipeline()) {
+ assertV7ForwardIndex(directory, segmentReader, rawForwardConfig,
column);
+ } else {
+ assertLegacyForwardIndex(directory, segmentReader,
rawForwardConfig, column);
+ }
+ }
+ }
+ }
+ }
+
+ private static void assertV7ForwardIndex(SegmentDirectory directory,
SegmentDirectory.Reader segmentReader,
+ FieldIndexConfigs rawForwardConfig, MatrixColumn column)
+ throws Exception {
+ ColumnMetadata metadata = rawColumnMetadata(directory, column._column);
+ PinotDataBuffer forwardIndexBuffer =
segmentReader.getIndexFor(column._column, StandardIndexes.forward());
+ try (ForwardIndexReader<?> forwardReader =
StandardIndexes.forward().getReaderFactory()
+ .createIndexReader(segmentReader, rawForwardConfig, metadata)) {
+ assertTrue(forwardReader instanceof FixedByteChunkSVForwardIndexReaderV7,
+ column + " should use the V7 codec-pipeline reader, got " +
forwardReader.getClass());
+
assertEquals(FixedByteChunkSVForwardIndexReaderV7.readCodecSpec(forwardIndexBuffer),
+ column._codecSpec._canonical, "Unexpected canonical spec in the V7
header of " + column);
+ assertNull(forwardReader.getCompressionType(), column + " must not
report a legacy ChunkCompressionType");
+ }
+ }
+
+ private static void assertLegacyForwardIndex(SegmentDirectory directory,
SegmentDirectory.Reader segmentReader,
+ FieldIndexConfigs rawForwardConfig, MatrixColumn column)
+ throws Exception {
+ ColumnMetadata metadata = rawColumnMetadata(directory, column._column);
+ try (ForwardIndexReader<?> forwardReader =
StandardIndexes.forward().getReaderFactory()
+ .createIndexReader(segmentReader, rawForwardConfig, metadata)) {
+ assertFalse(forwardReader instanceof
FixedByteChunkSVForwardIndexReaderV7,
+ column + " has no codecSpec and must not be routed to V7");
+ assertEquals(forwardReader.getCompressionType(),
expectedChunkCompressionType(column._compressionCodec),
+ "Unexpected ChunkCompressionType for " + column);
+ }
+ }
+
+ /// The [ChunkCompressionType] a legacy [CompressionCodec] maps to on a raw
index. `LZ4` is the only
+ /// one that is not name-identical: the var-byte V4 writer upgrades it to
`LZ4_LENGTH_PREFIXED` and
+ /// its reader maps that back to `LZ4`, so the expectation is `LZ4` for
every stored type.
+ private static ChunkCompressionType
expectedChunkCompressionType(CompressionCodec codec) {
+ switch (codec) {
+ case PASS_THROUGH:
+ return ChunkCompressionType.PASS_THROUGH;
+ case SNAPPY:
+ return ChunkCompressionType.SNAPPY;
+ case ZSTANDARD:
+ return ChunkCompressionType.ZSTANDARD;
+ case LZ4:
+ return ChunkCompressionType.LZ4;
+ case GZIP:
+ return ChunkCompressionType.GZIP;
+ default:
+ throw new IllegalStateException("Unexpected raw compression codec in
the matrix: " + codec);
+ }
+ }
+
+ private static ColumnMetadata rawColumnMetadata(SegmentDirectory directory,
String column) {
+ ColumnMetadata metadata =
directory.getSegmentMetadata().getColumnMetadataFor(column);
+ assertNotNull(metadata, "Missing metadata for " + column + " in " +
directory.getPath());
+ assertEquals(metadata.getForwardIndexEncoding(),
FieldConfig.EncodingType.RAW);
+ assertFalse(metadata.hasDictionary());
+ return metadata;
+ }
+
+ /// The table config as the controller stored it — not the one built in this
test — is what segment
+ /// generation reads, so every variant in the matrix must still be there
after the round trip.
+ @Test
+ public void testStoredTableConfigCarriesEveryVariant() {
+ // Read through the shared suite's controller: in suite mode this instance
never started one.
+ TableConfig storedTableConfig =
getSharedHelixResourceManager().getOfflineTableConfig(getTableName());
+ assertNotNull(storedTableConfig, "Controller has no stored OFFLINE table
config for " + getTableName());
+ List<String> noDictionaryColumns =
storedTableConfig.getIndexingConfig().getNoDictionaryColumns();
+ assertNotNull(noDictionaryColumns, "Stored table config lost
noDictionaryColumns");
+
+ Map<String, FieldConfig> storedFieldConfigs = new LinkedHashMap<>();
+ for (FieldConfig fieldConfig : storedTableConfig.getFieldConfigList()) {
+ storedFieldConfigs.put(fieldConfig.getName(), fieldConfig);
+ }
+ assertEquals(storedFieldConfigs.size(), MATRIX.size() + 1,
+ "Unexpected number of stored FieldConfigs: " +
storedFieldConfigs.keySet());
+
+ for (MatrixColumn column : MATRIX) {
+ FieldConfig fieldConfig = storedFieldConfigs.get(column._column);
+ assertNotNull(fieldConfig, "Stored table config has no FieldConfig for "
+ column);
+ assertEquals(fieldConfig.getEncodingType(),
FieldConfig.EncodingType.RAW, column + " must stay RAW");
+ assertTrue(noDictionaryColumns.contains(column._column),
+ "Expected " + column._column + " in the stored noDictionaryColumns");
+ JsonNode indexes = fieldConfig.getIndexes();
+ JsonNode forward = indexes == null ? null : indexes.get("forward");
+ assertNotNull(forward, "Stored FieldConfig for " + column + " lost its
indexes.forward block");
+ if (column.isCodecPipeline()) {
+ assertTrue(forward.hasNonNull("codecSpec"), "Stored FieldConfig for "
+ column + " lost its codecSpec");
+ assertEquals(forward.get("codecSpec").asText(),
column._codecSpec._configured,
+ "codecSpec for " + column + " did not survive the table config
round trip");
+ assertNull(fieldConfig.getCompressionCodec(),
+ column + " must not gain a compressionCodec; it is mutually
exclusive with codecSpec");
+ } else {
+ assertFalse(forward.hasNonNull("codecSpec"), column + " must not gain
a codecSpec");
+ assertEquals(fieldConfig.getCompressionCodec(),
column._compressionCodec,
+ "compressionCodec for " + column + " did not survive the table
config round trip");
+ }
+ }
+
+ FieldConfig dictionaryFieldConfig = storedFieldConfigs.get(DICT_STR_COL);
+ assertNotNull(dictionaryFieldConfig, "Stored table config has no
FieldConfig for " + DICT_STR_COL);
+ assertEquals(dictionaryFieldConfig.getEncodingType(),
FieldConfig.EncodingType.DICTIONARY);
+ assertFalse(noDictionaryColumns.contains(DICT_STR_COL),
+ DICT_STR_COL + " must not be in the stored noDictionaryColumns");
+ }
+}
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineRealtimeIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineRealtimeIntegrationTest.java
new file mode 100644
index 00000000000..53456b1ea8c
--- /dev/null
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/custom/CodecPipelineRealtimeIntegrationTest.java
@@ -0,0 +1,432 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.integration.tests.custom;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.node.ObjectNode;
+import java.io.File;
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.function.Consumer;
+import javax.annotation.Nullable;
+import org.apache.avro.Schema.Field;
+import org.apache.avro.Schema.Type;
+import org.apache.avro.file.DataFileWriter;
+import org.apache.avro.generic.GenericData;
+import org.apache.pinot.segment.local.data.manager.SegmentDataManager;
+import org.apache.pinot.segment.local.data.manager.TableDataManager;
+import
org.apache.pinot.segment.local.indexsegment.immutable.ImmutableSegmentImpl;
+import
org.apache.pinot.segment.local.segment.index.readers.forward.FixedByteChunkSVForwardIndexReaderV7;
+import org.apache.pinot.segment.local.segment.store.SegmentLocalFSDirectory;
+import org.apache.pinot.segment.spi.IndexSegment;
+import org.apache.pinot.segment.spi.compression.ChunkCompressionType;
+import org.apache.pinot.segment.spi.index.StandardIndexes;
+import org.apache.pinot.segment.spi.index.reader.ForwardIndexReader;
+import org.apache.pinot.segment.spi.memory.PinotDataBuffer;
+import org.apache.pinot.segment.spi.store.SegmentDirectory;
+import org.apache.pinot.server.starter.helix.BaseServerStarter;
+import org.apache.pinot.spi.config.table.FieldConfig;
+import org.apache.pinot.spi.config.table.FieldConfig.CompressionCodec;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.utils.CommonConstants;
+import org.apache.pinot.spi.utils.JsonUtils;
+import org.apache.pinot.spi.utils.ReadMode;
+import org.apache.pinot.spi.utils.builder.TableNameBuilder;
+import org.apache.pinot.util.TestUtils;
+import org.testng.annotations.BeforeClass;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertTrue;
+
+
+/// Realtime coverage for raw forward-index encodings, limited to what is
realtime-specific: the
+/// configured codec is ignored while rows are consuming and is applied when a
consuming segment is
+/// committed and converted to an immutable segment.
+///
+/// This deliberately does NOT extend [CodecPipelineIntegrationTest] and
deliberately does not repeat
+/// its matrix. Once a segment exists on disk, the codec matrix, both query
engines, dictionary
+/// coexistence and the chunk-boundary read paths are all table-type agnostic,
so inheriting them
+/// re-ran the whole offline matrix against a realtime table for no added
coverage. Three columns are
+/// enough here, because the commit path is what is under test rather than the
codecs:
+///
+/// - `longSpecDeltaLz4` (`DELTA,LZ4`) and `intSpecT64Zstd` (`T64,ZSTD(3)`) —
one LONG and one INT
+/// `codecSpec` column, so both the long and the int side of the V7 writer
run through a commit.
+/// - `longLegacyPassThrough` — legacy `compressionCodec: PASS_THROUGH`, the
uncompressed baseline and
+/// the negative control for the routing branch in
`ForwardIndexCreatorFactory`: only `codecSpec`
+/// selects the V7 writer, so a regression that routed every raw column to
V7 would still pass every
+/// other assertion here.
+///
+/// `ForwardIndexType.createMutableIndex` never looks at `codecSpec` — a
no-dictionary INT/LONG column
+/// always gets `FixedByteSVMutableForwardIndex` — so consuming rows are
uncompressed by design. The
+/// inherited `setUp` waits for every row to be queryable before this class
force-commits, so that
+/// consuming path is already exercised by setup. The flush threshold is below
the row count, so
+/// segments seal both by threshold and by `forceCommit`.
+@Test(suiteName = "CustomClusterIntegrationTest")
+public class CodecPipelineRealtimeIntegrationTest extends
CustomDataQueryClusterIntegrationTest {
+
+ private static final String TABLE_NAME =
"CodecPipelineRealtimeIntegrationTest";
+ private static final int NUM_DOCS = 600;
+ private static final int SEGMENT_FLUSH_SIZE = 250;
+ private static final int V7_TARGET_DOCS_PER_CHUNK = 64;
+ private static final long LONG_VALUE_SCALE = 1_000_000_000L;
+ private static final long FORCE_COMMIT_TIMEOUT_MS = 120_000L;
+ private static final long SEGMENT_LOAD_TIMEOUT_MS = 120_000L;
+ /// Segments sealed by the flush threshold, plus the one force-committed
remainder.
+ private static final int MIN_COMMITTED_SEGMENTS = NUM_DOCS /
SEGMENT_FLUSH_SIZE + 1;
+
+ private static final String TIME_COL = "ts";
+
+ /// The columns under test: two `codecSpec` shapes and the uncompressed
legacy control. Everything
+ /// below — schema, Avro fields, `noDictionaryColumns`, `FieldConfig`s, the
format assertions and the
+ /// queries — is derived from this list.
+ private static final List<RealtimeColumn> COLUMNS = List.of(
+ new RealtimeColumn("longSpecDeltaLz4", DataType.LONG, "DELTA,LZ4",
"DELTA,LZ4", null),
+ new RealtimeColumn("intSpecT64Zstd", DataType.INT, "T64,ZSTD",
"T64,ZSTD(3)", null),
+ new RealtimeColumn("longLegacyPassThrough", DataType.LONG, null, null,
CompressionCodec.PASS_THROUGH));
+
+ /// Doc ids at and around a V7 chunk boundary (63/64), a segment boundary
(249/250), and the last row
+ /// of the final force-committed segment.
+ private static final int[] POINT_LOOKUP_IDS = {0, 63, 64, 249, 250, 599};
+ private static final String POINT_LOOKUP_ID_LIST = "0, 63, 64, 249, 250,
599";
+
+ /// One column under test. Exactly one of the `codecSpec` pair and
[#_compressionCodec] is set,
+ /// because `ForwardIndexConfig` rejects both at once.
+ private static final class RealtimeColumn {
+ final String _column;
+ final DataType _dataType;
+ @Nullable
+ final String _codecSpec;
+ /// The canonical spec the codec runtime freezes into the V7 header;
`T64,ZSTD` materializes ZSTD's
+ /// default level, so it is not always the configured string.
+ @Nullable
+ final String _canonicalCodecSpec;
+ @Nullable
+ final CompressionCodec _compressionCodec;
+
+ private RealtimeColumn(String column, DataType dataType, @Nullable String
codecSpec,
+ @Nullable String canonicalCodecSpec, @Nullable CompressionCodec
compressionCodec) {
+ _column = column;
+ _dataType = dataType;
+ _codecSpec = codecSpec;
+ _canonicalCodecSpec = canonicalCodecSpec;
+ _compressionCodec = compressionCodec;
+ }
+
+ boolean isCodecPipeline() {
+ return _codecSpec != null;
+ }
+
+ @Override
+ public String toString() {
+ return _column + "[" + _dataType + ", " + (_codecSpec != null ?
"codecSpec=" + _codecSpec
+ : "compressionCodec=" + _compressionCodec) + "]";
+ }
+ }
+
+ /// INT columns hold the doc id; LONG columns scale it past the INT range so
a LONG column cannot
+ /// accidentally pass with 32-bit decoding. Both are a pure function of the
doc id.
+ private static long valueFor(DataType dataType, long docId) {
+ return dataType == DataType.INT ? docId : docId * LONG_VALUE_SCALE;
+ }
+
+ /// No sorted column: the forward index is the only thing under test, and
the base class would
+ /// otherwise make the time column sorted, which reorders rows at commit.
+ @Override
+ protected String getSortedColumn() {
+ return null;
+ }
+
+ @Override
+ public String getTableName() {
+ return TABLE_NAME;
+ }
+
+ @Override
+ public boolean isRealtimeTable() {
+ return true;
+ }
+
+ @Override
+ protected int getNumKafkaPartitions() {
+ return 1;
+ }
+
+ @Override
+ public int getNumAvroFiles() {
+ // One file and one partition, so rows reach the single consumer in `ts`
order. That makes the
+ // segment and chunk boundaries the point lookups below aim at
deterministic.
+ return 1;
+ }
+
+ @Override
+ protected int getRealtimeSegmentFlushSize() {
+ // < NUM_DOCS so a segment seals during consumption; forceCommit seals the
remainder.
+ return SEGMENT_FLUSH_SIZE;
+ }
+
+ @Override
+ public String getTimeColumnName() {
+ return TIME_COL;
+ }
+
+ @Override
+ protected long getCountStarResult() {
+ return NUM_DOCS;
+ }
+
+ @Override
+ public Schema createSchema() {
+ Schema.SchemaBuilder builder = new
Schema.SchemaBuilder().setSchemaName(getTableName());
+ for (RealtimeColumn column : COLUMNS) {
+ if (column.isCodecPipeline()) {
+ builder.addMetric(column._column, column._dataType);
+ } else {
+ // Deliberately a dimension, not a metric: metric columns already
default to PASS_THROUGH, so
+ // asserting PASS_THROUGH on a metric would pass even if the explicit
compressionCodec were
+ // dropped. Dimensions default to LZ4, which makes the negative
control meaningful.
+ builder.addSingleValueDimension(column._column, column._dataType);
+ }
+ }
+ builder.addDateTimeField(TIME_COL, DataType.LONG, "1:MILLISECONDS:EPOCH",
"1:MILLISECONDS");
+ return builder.build();
+ }
+
+ @Override
+ public List<File> createAvroFiles()
+ throws IOException {
+ org.apache.avro.Schema avroSchema =
org.apache.avro.Schema.createRecord("codecRealtimeRecord", null, null, false);
+ List<Field> fields = new ArrayList<>();
+ for (RealtimeColumn column : COLUMNS) {
+ Type avroType = column._dataType == DataType.INT ? Type.INT : Type.LONG;
+ fields.add(new Field(column._column,
org.apache.avro.Schema.create(avroType), null, null));
+ }
+ fields.add(new Field(TIME_COL, org.apache.avro.Schema.create(Type.LONG),
null, null));
+ avroSchema.setFields(fields);
+
+ try (AvroFilesAndWriters avroFilesAndWriters =
createAvroFilesAndWriters(avroSchema)) {
+ List<DataFileWriter<GenericData.Record>> writers =
avroFilesAndWriters.getWriters();
+ for (int docId = 0; docId < NUM_DOCS; docId++) {
+ GenericData.Record record = new GenericData.Record(avroSchema);
+ for (RealtimeColumn column : COLUMNS) {
+ if (column._dataType == DataType.INT) {
+ record.put(column._column, docId);
+ } else {
+ record.put(column._column, valueFor(DataType.LONG, docId));
+ }
+ }
+ record.put(TIME_COL, (long) docId);
+ writers.get(docId % getNumAvroFiles()).append(record);
+ }
+ return avroFilesAndWriters.getAvroFiles();
+ }
+ }
+
+ @Override
+ protected List<String> getNoDictionaryColumns() {
+ List<String> noDictionaryColumns = new ArrayList<>(COLUMNS.size());
+ for (RealtimeColumn column : COLUMNS) {
+ noDictionaryColumns.add(column._column);
+ }
+ return noDictionaryColumns;
+ }
+
+ @Override
+ protected List<FieldConfig> getFieldConfigs() {
+ List<FieldConfig> fieldConfigs = new ArrayList<>(COLUMNS.size());
+ for (RealtimeColumn column : COLUMNS) {
+ ObjectNode forward = JsonUtils.newObjectNode();
+ // Keep chunks well below the flush size so committed segments hold
several chunks.
+ forward.put("targetDocsPerChunk", V7_TARGET_DOCS_PER_CHUNK);
+ if (column.isCodecPipeline()) {
+ forward.put("codecSpec", column._codecSpec);
+ }
+ ObjectNode indexes = JsonUtils.newObjectNode();
+ indexes.set("forward", forward);
+ FieldConfig.Builder builder = new FieldConfig.Builder(column._column)
+ .withEncodingType(FieldConfig.EncodingType.RAW)
+ .withIndexes(indexes);
+ if (!column.isCodecPipeline()) {
+ builder.withCompressionCodec(column._compressionCodec);
+ }
+ fieldConfigs.add(builder.build());
+ }
+ return fieldConfigs;
+ }
+
+ @Override
+ @BeforeClass
+ public void setUp()
+ throws Exception {
+ // Loads the schema and realtime table config, pushes all rows into Kafka,
and waits until every
+ // row is queryable — which happens while the last segment is still
consuming.
+ super.setUp();
+ forceCommitAndWait();
+ // The commit job completing means the segments are committed in
ZooKeeper; the servers still
+ // have to load them before their forward-index format can be inspected or
queried.
+ TestUtils.waitForCondition(aVoid -> countCommittedSegments() >=
MIN_COMMITTED_SEGMENTS, 1_000L,
+ SEGMENT_LOAD_TIMEOUT_MS, "Timed out waiting for " +
MIN_COMMITTED_SEGMENTS
+ + " committed realtime segments to load");
+ }
+
+ /// The realtime-specific assertion: the configured forward-index format is
applied when a consuming
+ /// segment is committed and converted to an immutable segment. A
`codecSpec` column becomes a
+ /// self-describing V7 index carrying the canonical spec and no legacy
ChunkCompressionType, while
+ /// the legacy column in the same segment stays on the legacy writer and
reports `PASS_THROUGH`.
+ @Test
+ public void testCommittedSegmentsUseConfiguredForwardIndexFormats() {
+ // A lower bound rather than an equality: every committed segment found is
asserted, so an extra
+ // seal is not a failure, but finding fewer than the guaranteed ones is.
+ int inspected =
forEachCommittedSegment(CodecPipelineRealtimeIntegrationTest::assertForwardIndexFormats);
+ assertTrue(inspected >= MIN_COMMITTED_SEGMENTS,
+ "Expected at least " + MIN_COMMITTED_SEGMENTS + " committed realtime
segments, inspected " + inspected);
+ }
+
+ private int countCommittedSegments() {
+ return forEachCommittedSegment(segment -> {
+ });
+ }
+
+ /// Visits every committed (immutable) segment of this table loaded on the
shared servers and
+ /// returns how many there were. Consuming segments are skipped: while
consuming, the rows live in
+ /// a mutable forward index that ignores `codecSpec` by design.
+ private int forEachCommittedSegment(Consumer<ImmutableSegmentImpl> visitor) {
+ String realtimeTableName =
TableNameBuilder.REALTIME.tableNameWithType(getTableName());
+ int visited = 0;
+ for (BaseServerStarter serverStarter : getSharedServerStarters()) {
+ TableDataManager tableDataManager =
serverStarter.getServerInstance().getInstanceDataManager()
+ .getTableDataManager(realtimeTableName);
+ if (tableDataManager == null) {
+ continue;
+ }
+ List<SegmentDataManager> segmentDataManagers =
tableDataManager.acquireAllSegments();
+ try {
+ for (SegmentDataManager segmentDataManager : segmentDataManagers) {
+ IndexSegment segment = segmentDataManager.getSegment();
+ if (segment instanceof ImmutableSegmentImpl) {
+ visitor.accept((ImmutableSegmentImpl) segment);
+ visited++;
+ }
+ }
+ } finally {
+ for (SegmentDataManager segmentDataManager : segmentDataManagers) {
+ tableDataManager.releaseSegment(segmentDataManager);
+ }
+ }
+ }
+ return visited;
+ }
+
+ /// End-to-end read back of the committed segments through both query
engines: an aggregate over
+ /// every row, the whole-table cross-codec equality check, and point lookups
at V7 chunk and segment
+ /// boundaries — all derived from [#COLUMNS] so they cannot drift from it.
+ @Test(dataProvider = "useBothQueryEngines")
+ public void testQueryCommittedSegments(boolean useMultiStageQueryEngine)
+ throws Exception {
+ setUseMultiStageQueryEngine(useMultiStageQueryEngine);
+
+ List<String> sums = new ArrayList<>(COLUMNS.size());
+ List<String> selectList = new ArrayList<>(COLUMNS.size() + 1);
+ List<String> predicates = new ArrayList<>(COLUMNS.size());
+ selectList.add(TIME_COL);
+ for (RealtimeColumn column : COLUMNS) {
+ sums.add("SUM(" + column._column + ")");
+ selectList.add(column._column);
+ predicates.add(column._dataType == DataType.INT ? column._column + " = "
+ TIME_COL
+ : column._column + " = " + TIME_COL + " * " + LONG_VALUE_SCALE);
+ }
+
+ JsonNode sumResult = postQuery("SELECT " + String.join(", ", sums) + "
FROM " + getTableName());
+ JsonNode sumRow = sumResult.get("resultTable").get("rows").get(0);
+ for (int i = 0; i < COLUMNS.size(); i++) {
+ RealtimeColumn column = COLUMNS.get(i);
+ assertEquals(sumRow.get(i).asLong(), valueFor(column._dataType, (long)
NUM_DOCS * (NUM_DOCS - 1) / 2),
+ "Unexpected SUM for " + column + " after commit");
+ }
+
+ // Every row of every column agrees with `ts`, so a per-row decoding error
that happens to
+ // preserve the SUM cannot slip through.
+ JsonNode agreement = postQuery(
+ "SELECT COUNT(*) FROM " + getTableName() + " WHERE " + String.join("
AND ", predicates));
+
assertEquals(agreement.get("resultTable").get("rows").get(0).get(0).asLong(),
NUM_DOCS,
+ "Not every row agrees across the committed realtime columns: " +
String.join(" AND ", predicates));
+
+ JsonNode result = postQuery("SELECT " + String.join(", ", selectList) + "
FROM " + getTableName() + " WHERE "
+ + TIME_COL + " IN (" + POINT_LOOKUP_ID_LIST + ") ORDER BY " +
TIME_COL);
+ JsonNode rows = result.get("resultTable").get("rows");
+ assertEquals(rows.size(), POINT_LOOKUP_IDS.length, "Unexpected
point-lookup row count after commit");
+ for (int rowId = 0; rowId < POINT_LOOKUP_IDS.length; rowId++) {
+ int docId = POINT_LOOKUP_IDS[rowId];
+ JsonNode row = rows.get(rowId);
+ assertEquals(row.get(0).asInt(), docId, "Unexpected point-lookup order
after commit");
+ for (int i = 0; i < COLUMNS.size(); i++) {
+ RealtimeColumn column = COLUMNS.get(i);
+ assertEquals(row.get(i + 1).asLong(), valueFor(column._dataType,
docId),
+ "Wrong " + column + " for ts=" + docId);
+ }
+ }
+ }
+
+ private static void assertForwardIndexFormats(ImmutableSegmentImpl segment) {
+ try (SegmentDirectory directory = new
SegmentLocalFSDirectory(segment.getSegmentMetadata().getIndexDir(),
+ ReadMode.mmap); SegmentDirectory.Reader segmentReader =
directory.createReader()) {
+ for (RealtimeColumn column : COLUMNS) {
+ ForwardIndexReader<?> reader =
segment.getDataSource(column._column).getForwardIndex();
+ if (!column.isCodecPipeline()) {
+ assertFalse(reader instanceof FixedByteChunkSVForwardIndexReaderV7,
+ column + " has no codecSpec and must not be routed to V7");
+ assertEquals(reader.getCompressionType(),
ChunkCompressionType.PASS_THROUGH,
+ "Unexpected ChunkCompressionType for " + column);
+ continue;
+ }
+ assertTrue(reader instanceof FixedByteChunkSVForwardIndexReaderV7,
+ column + " was routed to " + reader.getClass().getSimpleName());
+ PinotDataBuffer forwardIndexBuffer =
segmentReader.getIndexFor(column._column, StandardIndexes.forward());
+
assertEquals(FixedByteChunkSVForwardIndexReaderV7.readCodecSpec(forwardIndexBuffer),
+ column._canonicalCodecSpec, "Unexpected canonical spec in the V7
header of " + column);
+ assertNull(reader.getCompressionType(), column + " must not report a
legacy ChunkCompressionType");
+ }
+ } catch (Exception e) {
+ throw new AssertionError("Failed to inspect the forward-index formats of
" + segment.getSegmentName(), e);
+ }
+ }
+
+ private void forceCommitAndWait()
+ throws Exception {
+ String realtimeTableName =
TableNameBuilder.REALTIME.tableNameWithType(getTableName());
+ String response =
getOrCreateAdminClient().getTableClient().forceCommit(realtimeTableName);
+ String jobId =
JsonUtils.stringToJsonNode(response).get("forceCommitJobId").asText();
+ TestUtils.waitForCondition(aVoid -> isForceCommitComplete(jobId), 1_000L,
FORCE_COMMIT_TIMEOUT_MS,
+ "Timed out waiting for forceCommit job: " + jobId);
+ }
+
+ private boolean isForceCommitComplete(String jobId) {
+ try {
+ String response =
getOrCreateAdminClient().getTableClient().getForceCommitJobStatus(jobId);
+ JsonNode status = JsonUtils.stringToJsonNode(response);
+ return
status.get(CommonConstants.ControllerJob.NUM_CONSUMING_SEGMENTS_YET_TO_BE_COMMITTED).asInt(-1)
== 0;
+ } catch (Exception e) {
+ return false;
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]