This is an automated email from the ASF dual-hosted git repository.

KKcorps 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 816375d50c7 Take upsert/dedup primary key columns from the metadata 
manager context (#19726)
816375d50c7 is described below

commit 816375d50c7c8cc8a34895f0862e2327ab0de489
Author: Chaitanya Deepthi <[email protected]>
AuthorDate: Wed Sep 30 20:06:18 2026 -0700

    Take upsert/dedup primary key columns from the metadata manager context 
(#19726)
---
 .../indexsegment/mutable/MutableSegmentImpl.java   | 14 +++++-
 .../mutable/MutableSegmentDedupTest.java           | 44 +++++++++++++++++
 .../mutable/MutableSegmentImplUpsertTest.java      | 55 ++++++++++++++++++++++
 3 files changed, 111 insertions(+), 2 deletions(-)

diff --git 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImpl.java
 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImpl.java
index 0db9174f619..aa699f38570 100644
--- 
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImpl.java
+++ 
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImpl.java
@@ -176,8 +176,10 @@ public class MutableSegmentImpl implements MutableSegment {
   private final Collection<ComplexFieldSpec> _physicalComplexFieldSpecs;
   private final PartitionDedupMetadataManager _partitionDedupMetadataManager;
   private final String _dedupTimeColumn;
+  private final List<String> _dedupPrimaryKeyColumns;
   private final PartitionUpsertMetadataManager _partitionUpsertMetadataManager;
   private final boolean _isPartialUpsert;
+  private final List<String> _upsertPrimaryKeyColumns;
   private final List<String> _upsertComparisonColumns;
   private final String _deleteRecordColumn;
   private final boolean _upsertDropOutOfOrderRecord;
@@ -457,6 +459,9 @@ public class MutableSegmentImpl implements MutableSegment {
     _dedupTimeColumn =
         _partitionDedupMetadataManager != null ? 
_partitionDedupMetadataManager.getContext().getDedupTimeColumn()
             : null;
+    _dedupPrimaryKeyColumns =
+        _partitionDedupMetadataManager != null ? 
_partitionDedupMetadataManager.getContext().getPrimaryKeyColumns()
+            : null;
 
     _partitionUpsertMetadataManager = 
config.getPartitionUpsertMetadataManager();
     if (_partitionUpsertMetadataManager != null) {
@@ -464,6 +469,7 @@ public class MutableSegmentImpl implements MutableSegment {
           "Metrics aggregation and upsert cannot be enabled together");
       UpsertContext upsertContext = 
_partitionUpsertMetadataManager.getContext();
       _isPartialUpsert = upsertContext.getUpsertMode() == 
UpsertConfig.Mode.PARTIAL;
+      _upsertPrimaryKeyColumns = upsertContext.getPrimaryKeyColumns();
       _upsertComparisonColumns = upsertContext.getComparisonColumns();
       _deleteRecordColumn = upsertContext.getDeleteRecordColumn();
       _upsertDropOutOfOrderRecord = upsertContext.isDropOutOfOrderRecord();
@@ -477,6 +483,7 @@ public class MutableSegmentImpl implements MutableSegment {
       }
     } else {
       _isPartialUpsert = false;
+      _upsertPrimaryKeyColumns = null;
       _upsertComparisonColumns = null;
       _deleteRecordColumn = null;
       _upsertDropOutOfOrderRecord = false;
@@ -778,7 +785,7 @@ public class MutableSegmentImpl implements MutableSegment {
   }
 
   private DedupRecordInfo getDedupRecordInfo(GenericRow row) {
-    PrimaryKey primaryKey = row.getPrimaryKey(_schema.getPrimaryKeyColumns());
+    PrimaryKey primaryKey = row.getPrimaryKey(_dedupPrimaryKeyColumns);
     // it is okay not having dedup time column if metadata ttl is not enabled
     if (_dedupTimeColumn == null) {
       return new DedupRecordInfo(primaryKey);
@@ -788,7 +795,10 @@ public class MutableSegmentImpl implements MutableSegment {
   }
 
   private RecordInfo getRecordInfo(GenericRow row, int docId) {
-    PrimaryKey primaryKey = row.getPrimaryKey(_schema.getPrimaryKeyColumns());
+    // Take the key from the metadata manager's context, not this segment's 
schema: the context is fixed for the
+    // life of the manager, while a new consuming segment picks up the latest 
schema. Reading the schema here lets
+    // a primary key change desync the two and hash the same row under two 
different keys.
+    PrimaryKey primaryKey = row.getPrimaryKey(_upsertPrimaryKeyColumns);
     Comparable comparisonValue = getComparisonValue(row);
     boolean deleteRecord = _deleteRecordColumn != null && 
BooleanUtils.toBoolean(row.getValue(_deleteRecordColumn));
     return new RecordInfo(primaryKey, docId, comparisonValue, deleteRecord);
diff --git 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentDedupTest.java
 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentDedupTest.java
index 2eb642305fa..d24059be4d6 100644
--- 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentDedupTest.java
+++ 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentDedupTest.java
@@ -58,6 +58,8 @@ public class MutableSegmentDedupTest implements 
PinotBuffersAfterMethodCheckRule
   private static final String DATA_FILE_PATH = "data/test_dedup_data.json";
   private static final String RAW_TABLE_NAME = "testTable";
   private static final String TIME_COLUMN = "secondsSinceEpoch";
+  private static final String PRIMARY_KEY_COLUMN = "event_id";
+  private static final String DEDUP_TIME_COLUMN = "dedupTime";
 
   private MutableSegmentImpl _mutableSegmentImpl;
 
@@ -147,6 +149,48 @@ public class MutableSegmentDedupTest implements 
PinotBuffersAfterMethodCheckRule
     }
   }
 
+  /// The dedup metadata manager captures the primary key columns once, when 
the server starts, while a consuming
+  /// segment created later picks up whatever schema is current. Both must key 
records identically, otherwise the
+  /// same record hashes under two different keys and duplicates go undetected.
+  @Test
+  public void testPrimaryKeyColumnsTakenFromDedupContext()
+      throws Exception {
+    FileUtils.forceMkdir(TEMP_DIR);
+    URL schemaResourceUrl = 
this.getClass().getClassLoader().getResource(SCHEMA_FILE_PATH);
+    Assert.assertNotNull(schemaResourceUrl);
+    // Schema as it was when the metadata manager was created: a single 
primary key column.
+    Schema managerSchema = Schema.fromFile(new 
File(schemaResourceUrl.getFile()));
+    // Schema the consuming segment was created with: a second column has 
since joined the primary key.
+    Schema segmentSchema = Schema.fromFile(new 
File(schemaResourceUrl.getFile()));
+    segmentSchema.setPrimaryKeyColumns(List.of(PRIMARY_KEY_COLUMN, 
DEDUP_TIME_COLUMN));
+
+    PartitionDedupMetadataManager partitionDedupMetadataManager =
+        getTableDedupMetadataManager(managerSchema, new 
DedupConfig()).getOrCreatePartitionManager(0);
+    try {
+      _mutableSegmentImpl = 
MutableSegmentImplTestUtils.createMutableSegmentImpl(segmentSchema, true, 
TIME_COLUMN, null,
+          partitionDedupMetadataManager);
+      // Two records sharing event_id, differing only on the column that the 
segment schema treats as part of the
+      // key but the metadata manager does not.
+      _mutableSegmentImpl.index(createRow("aa", 1L, 1567205396L), null);
+      _mutableSegmentImpl.index(createRow("aa", 2L, 1567205397L), null);
+
+      // Keyed on the manager's columns these are the same record, so the 
second must have been deduped away.
+      Assert.assertEquals(_mutableSegmentImpl.getNumDocsIndexed(), 1);
+    } finally {
+      partitionDedupMetadataManager.stop();
+      partitionDedupMetadataManager.close();
+    }
+  }
+
+  private static GenericRow createRow(String eventId, long dedupTime, long 
timeValue) {
+    GenericRow row = new GenericRow();
+    row.putValue(PRIMARY_KEY_COLUMN, eventId);
+    row.putValue("description", "d");
+    row.putValue(DEDUP_TIME_COLUMN, dedupTime);
+    row.putValue(TIME_COLUMN, timeValue);
+    return row;
+  }
+
   @Test
   public void testDedupDisabled()
       throws Exception {
diff --git 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImplUpsertTest.java
 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImplUpsertTest.java
index d96da49e0e8..2140e7e8d27 100644
--- 
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImplUpsertTest.java
+++ 
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/indexsegment/mutable/MutableSegmentImplUpsertTest.java
@@ -54,6 +54,7 @@ public class MutableSegmentImplUpsertTest {
   private static final String RAW_TABLE_NAME = "testTable";
   private static final String TIME_COLUMN = "secondsSinceEpoch";
   private static final String OTHER_COMPARISON_COLUMN = 
"otherComparisonColumn";
+  private static final String PRIMARY_KEY_COLUMN = "event_id";
 
   private MutableSegmentImpl _mutableSegmentImpl;
   private PartitionUpsertMetadataManager _partitionUpsertMetadataManager;
@@ -128,6 +129,60 @@ public class MutableSegmentImplUpsertTest {
     testUpsertIngestion(createFullUpsertConfig(HashFunction.MURMUR3));
   }
 
+  /// The upsert metadata manager captures the primary key columns once, when 
the server starts. A consuming
+  /// segment created later picks up whatever schema is current, which may 
list a different set. Both must key
+  /// records identically, otherwise the same row hashes under two different 
keys and neither invalidates the other.
+  @Test
+  public void testPrimaryKeyColumnsTakenFromUpsertContext()
+      throws Exception {
+    TableConfig tableConfig = new 
TableConfigBuilder(TableType.REALTIME).setTableName(RAW_TABLE_NAME)
+        .setTimeColumnName(TIME_COLUMN)
+        .setUpsertConfig(createFullUpsertConfig(HashFunction.MURMUR3))
+        .setNullHandlingEnabled(true)
+        .build();
+    URL schemaResourceUrl = 
getClass().getClassLoader().getResource(SCHEMA_FILE_PATH);
+    assertNotNull(schemaResourceUrl);
+    // Schema as it was when the metadata manager was created: a single 
primary key column.
+    Schema managerSchema = Schema.fromFile(new 
File(schemaResourceUrl.getFile()));
+    // Schema the consuming segment was created with: a second column has 
since joined the primary key.
+    Schema segmentSchema = Schema.fromFile(new 
File(schemaResourceUrl.getFile()));
+    segmentSchema.setPrimaryKeyColumns(List.of(PRIMARY_KEY_COLUMN, 
OTHER_COMPARISON_COLUMN));
+
+    TableUpsertMetadataManager tableUpsertMetadataManager =
+        TableUpsertMetadataManagerFactory.create(new PinotConfiguration(), 
tableConfig, managerSchema,
+            mock(TableDataManager.class), null);
+    PartitionUpsertMetadataManager partitionUpsertMetadataManager =
+        tableUpsertMetadataManager.getOrCreatePartitionManager(0);
+    MutableSegmentImpl mutableSegment =
+        MutableSegmentImplTestUtils.createMutableSegmentImpl(segmentSchema, 
true, TIME_COLUMN,
+            partitionUpsertMetadataManager, null);
+    try {
+      // Two records sharing event_id, differing only on the column that the 
segment schema treats as part of the
+      // key but the metadata manager does not.
+      mutableSegment.index(createRow("aa", 1L, 1567205396L), null);
+      mutableSegment.index(createRow("aa", 2L, 1567205397L), null);
+
+      // Keyed on the manager's columns these are the same record, so the 
older doc must have been invalidated.
+      ImmutableRoaringBitmap bitmap = 
mutableSegment.getValidDocIds().getMutableRoaringBitmap();
+      assertEquals(bitmap.getCardinality(), 1);
+      assertFalse(bitmap.contains(0));
+      assertTrue(bitmap.contains(1));
+    } finally {
+      mutableSegment.destroy();
+      partitionUpsertMetadataManager.stop();
+      partitionUpsertMetadataManager.close();
+    }
+  }
+
+  private static GenericRow createRow(String eventId, long 
otherComparisonValue, long timeValue) {
+    GenericRow row = new GenericRow();
+    row.putValue(PRIMARY_KEY_COLUMN, eventId);
+    row.putValue("description", "d");
+    row.putValue(OTHER_COMPARISON_COLUMN, otherComparisonValue);
+    row.putValue(TIME_COLUMN, timeValue);
+    return row;
+  }
+
   @Test
   public void testMultipleComparisonColumns()
       throws Exception {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to