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]