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]

Reply via email to