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 03ed4257107 Require an explicit comparison column for offline upsert
tables (#19177)
03ed4257107 is described below
commit 03ed4257107a512314169b3967780fdab74df9f8
Author: Xiang Fu <[email protected]>
AuthorDate: Fri Aug 7 12:17:10 2026 -0700
Require an explicit comparison column for offline upsert tables (#19177)
Rolls back the offline upsert "no comparison column" support added in
#17789.
Offline upsert validation now matches realtime: a comparison column is
required, resolved from upsertConfig.comparisonColumns or the table's time
column. The implicit segment-creation-time fallback is removed as it caused
more problems than it solved. Offline upsert is not GA, so no
backward-compatibility shims are kept.
This is scoped to the comparison-column requirement only. The cross-replica
upsert lifecycle-consistency tie-break (getAuthoritativeUpdateOrCreationTime
and the ZK creation/push time plumbing) is left unchanged, as it is
orthogonal
and applies regardless of whether the comparison column is explicit.
- TableConfigUtils: offline upsert must configure a comparison column or a
time column (same resolution as realtime).
- BaseTableUpsertMetadataManager: revert comparison-column resolution to the
realtime behavior; fail fast with an actionable message instead of an NPE
when neither column is configured on the server load path.
- UpsertUtils: remove ConstantComparisonColumnReader and the constant-value
RecordInfoReader constructor (the segment-creation-time reader path).
- BasePartitionUpsertMetadataManager: remove the segment-creation-time
RecordInfoReader helper and revert the empty-comparison-column TTL guards.
Co-authored-by: Xiang Fu <[email protected]>
---
.../upsert/BasePartitionUpsertMetadataManager.java | 27 +++++++-----------
.../upsert/BaseTableUpsertMetadataManager.java | 11 ++++----
.../pinot/segment/local/upsert/UpsertUtils.java | 32 ----------------------
.../segment/local/utils/TableConfigUtils.java | 7 +++++
.../TableUpsertMetadataManagerFactoryTest.java | 19 +++++++++++++
.../segment/local/utils/TableConfigUtilsTest.java | 25 +++++++++++++++++
6 files changed, 66 insertions(+), 55 deletions(-)
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
index 4c370cfb2d3..166ff082500 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BasePartitionUpsertMetadataManager.java
@@ -312,16 +312,6 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
}
}
- /// Creates a RecordInfoReader for the given segment. When comparison
columns are configured, reads comparison values
- /// from the columns. When comparison columns are empty, uses segment
creation time as the comparison value.
- protected UpsertUtils.RecordInfoReader createRecordInfoReader(IndexSegment
segment) {
- if (_comparisonColumns.isEmpty()) {
- long segmentCreationTime = getAuthoritativeUpdateOrCreationTime(segment);
- return new UpsertUtils.RecordInfoReader(segment, _primaryKeyColumns,
segmentCreationTime, _deleteRecordColumn);
- }
- return new UpsertUtils.RecordInfoReader(segment, _primaryKeyColumns,
_comparisonColumns, _deleteRecordColumn);
- }
-
protected boolean isTTLEnabled() {
return _metadataTTL > 0 || _deletedKeysTTL > 0;
}
@@ -374,7 +364,7 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
protected void doAddSegment(ImmutableSegmentImpl segment) {
String segmentName = segment.getSegmentName();
_logger.info("Adding segment: {}, current primary key count: {}",
segmentName, getNumPrimaryKeys());
- if (isTTLEnabled() && !_comparisonColumns.isEmpty()) {
+ if (isTTLEnabled()) {
double maxComparisonValue = getMaxComparisonValue(segment);
_largestSeenComparisonValue.getAndUpdate(v -> Math.max(v,
maxComparisonValue));
if (isOutOfMetadataTTL(maxComparisonValue) &&
skipAddSegmentOutOfTTL(segment)) {
@@ -385,7 +375,8 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
if (!_enableSnapshot) {
deleteSnapshot(segment);
}
- try (UpsertUtils.RecordInfoReader recordInfoReader =
createRecordInfoReader(segment)) {
+ try (UpsertUtils.RecordInfoReader recordInfoReader = new
UpsertUtils.RecordInfoReader(segment, _primaryKeyColumns,
+ _comparisonColumns, _deleteRecordColumn)) {
Iterator<RecordInfo> recordInfoIterator =
UpsertUtils.getRecordInfoIterator(recordInfoReader,
segment.getSegmentMetadata().getTotalDocs());
addSegment(segment, null, null, recordInfoIterator);
@@ -451,14 +442,15 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
segment.enableUpsert(this, new ThreadSafeMutableRoaringBitmap(),
queryableDocIds);
return;
}
- if (isTTLEnabled() && !_comparisonColumns.isEmpty()) {
+ if (isTTLEnabled()) {
double maxComparisonValue = getMaxComparisonValue(segment);
_largestSeenComparisonValue.getAndUpdate(v -> Math.max(v,
maxComparisonValue));
if (isOutOfMetadataTTL(maxComparisonValue) &&
skipPreloadSegmentOutOfTTL(segment, validDocIds)) {
return;
}
}
- try (UpsertUtils.RecordInfoReader recordInfoReader =
createRecordInfoReader(segment)) {
+ try (UpsertUtils.RecordInfoReader recordInfoReader = new
UpsertUtils.RecordInfoReader(segment, _primaryKeyColumns,
+ _comparisonColumns, _deleteRecordColumn)) {
doPreloadSegment(segment, null, null,
UpsertUtils.getRecordInfoIterator(recordInfoReader, validDocIds));
} catch (Exception e) {
throw new RuntimeException(
@@ -604,13 +596,14 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
replaceSegment(segment, null, null, null, oldSegment);
return;
}
- if (isTTLEnabled() && !_comparisonColumns.isEmpty()) {
+ if (isTTLEnabled()) {
double maxComparisonValue = getMaxComparisonValue(segment);
_largestSeenComparisonValue.getAndUpdate(v -> Math.max(v,
maxComparisonValue));
// Segment might be uploaded directly to the table to replace an old
segment. So update the TTL watermark but
// we can't skip segment even if it's out of TTL as its validDocIds
bitmap is not updated yet.
}
- try (UpsertUtils.RecordInfoReader recordInfoReader =
createRecordInfoReader(segment)) {
+ try (UpsertUtils.RecordInfoReader recordInfoReader = new
UpsertUtils.RecordInfoReader(segment, _primaryKeyColumns,
+ _comparisonColumns, _deleteRecordColumn)) {
// Reload-only fast path for an upsert + TTL table. The incoming segment
carries a validDocIds snapshot ONLY when
// the reload flow placed it there (see
BaseTableDataManager.reloadSegment); segment commits and uploads always
// build a fresh segment without one and fall through to the full scan
below, unaffected. Rebuilding from just the
@@ -768,7 +761,7 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
try {
// Skip removing the upsert metadata of segment that is out of metadata
TTL. The expired metadata is removed
// while creating new consuming segment in batches.
- if (!_comparisonColumns.isEmpty() && isOutOfMetadataTTL(segment)) {
+ if (isOutOfMetadataTTL(segment)) {
_logger.info("Skip removing segment: {} because it's out of TTL",
segmentName);
} else {
doRemoveSegment(segment);
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BaseTableUpsertMetadataManager.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BaseTableUpsertMetadataManager.java
index 7c801e4634a..0c155fb308b 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BaseTableUpsertMetadataManager.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/BaseTableUpsertMetadataManager.java
@@ -60,13 +60,12 @@ public abstract class BaseTableUpsertMetadataManager
implements TableUpsertMetad
List<String> comparisonColumns = upsertConfig.getComparisonColumns();
if (comparisonColumns == null) {
+ // Fall back to the table's time column, same as realtime upsert.
Realtime always has a time column, but an
+ // offline upsert table might not, so fail fast with an actionable
message instead of an NPE.
String timeColumnName =
tableConfig.getValidationConfig().getTimeColumnName();
- if (timeColumnName != null) {
- comparisonColumns = List.of(timeColumnName);
- } else {
- // No comparison column and no time column: use segment creation time
for comparison
- comparisonColumns = List.of();
- }
+ Preconditions.checkState(timeColumnName != null,
+ "Upsert table: %s must have a comparison column or a time column
configured", _tableNameWithType);
+ comparisonColumns = List.of(timeColumnName);
}
// PartialUpsertHandler is not thread safe, so hand each partition a
factory rather than one shared instance.
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertUtils.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertUtils.java
index c2208ab43c6..494bae6a7a3 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertUtils.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertUtils.java
@@ -244,19 +244,6 @@ public class UpsertUtils {
}
}
- /// Constructor that uses a constant comparison value for all records.
- /// Used when no comparison columns are configured and segment creation
time is used as the comparison value.
- public RecordInfoReader(IndexSegment segment, List<String>
primaryKeyColumns,
- Comparable constantComparisonValue, @Nullable String
deleteRecordColumn) {
- _primaryKeyReader = new PrimaryKeyReader(segment, primaryKeyColumns);
- _comparisonColumnReader = new
ConstantComparisonColumnReader(constantComparisonValue);
- if (deleteRecordColumn != null) {
- _deleteRecordColumnReader = new PinotSegmentColumnReader(segment,
deleteRecordColumn);
- } else {
- _deleteRecordColumnReader = null;
- }
- }
-
public RecordInfo getRecordInfo(int docId) {
PrimaryKey primaryKey = _primaryKeyReader.getPrimaryKey(docId);
Comparable comparisonValue =
_comparisonColumnReader.getComparisonValue(docId);
@@ -338,23 +325,4 @@ public class UpsertUtils {
}
}
}
-
- /// A comparison column reader that returns a constant value for all records.
- /// Used when no comparison columns are configured and segment creation time
is used as the comparison value.
- public static class ConstantComparisonColumnReader implements
ComparisonColumnReader {
- private final Comparable _constantValue;
-
- public ConstantComparisonColumnReader(Comparable constantValue) {
- _constantValue = constantValue;
- }
-
- @Override
- public Comparable getComparisonValue(int docId) {
- return _constantValue;
- }
-
- @Override
- public void close() {
- }
- }
}
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
index 0e2868c795e..222a8396e78 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/utils/TableConfigUtils.java
@@ -1049,6 +1049,13 @@ public final class TableConfigUtils {
segmentPartitionConfig != null &&
MapUtils.isNotEmpty(segmentPartitionConfig.getColumnPartitionMap()),
"Offline upsert table must have segment partition config to ensure
correct partition-based "
+ "segment assignment. Configure segmentPartitionConfig in the
indexingConfig.");
+ // A comparison column is required, same as realtime upsert. It resolves
to the configured comparison columns, or
+ // the table's time column when comparison columns are not set. Unlike
realtime, offline tables do not require a
+ // time column, so reject the case where neither is configured (no
implicit segment creation time fallback).
+ Preconditions.checkState(
+ CollectionUtils.isNotEmpty(upsertConfig.getComparisonColumns())
+ || tableConfig.getValidationConfig().getTimeColumnName() != null,
+ "Offline upsert table must have a comparison column or a time column
configured");
}
if (upsertEnabled) {
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/TableUpsertMetadataManagerFactoryTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/TableUpsertMetadataManagerFactoryTest.java
index 4dfb2485029..0d062e1f8cb 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/TableUpsertMetadataManagerFactoryTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/TableUpsertMetadataManagerFactoryTest.java
@@ -39,6 +39,7 @@ import static org.mockito.Mockito.when;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.expectThrows;
public class TableUpsertMetadataManagerFactoryTest {
@@ -72,6 +73,24 @@ public class TableUpsertMetadataManagerFactoryTest {
instanceof ConcurrentMapPartitionUpsertMetadataManager);
}
+ @Test
+ public void testCreateFailsWithoutComparisonColumnOrTimeColumn() {
+ // An upsert table with neither a configured comparison column nor a time
column has no way to resolve a comparison
+ // column, so init() must fail fast rather than fall back to an implicit
value.
+ UpsertConfig upsertConfig = new UpsertConfig(UpsertConfig.Mode.FULL);
+ TableConfig tableConfig = new TableConfigBuilder(TableType.REALTIME)
+ .setTableName(RAW_TABLE_NAME)
+ .setUpsertConfig(upsertConfig)
+ .build();
+ TableDataManager tableDataManager = mock(TableDataManager.class);
+ when(tableDataManager.getTableDataDir()).thenReturn(new
File(RAW_TABLE_NAME));
+ IllegalStateException e = expectThrows(IllegalStateException.class,
+ () -> TableUpsertMetadataManagerFactory.create(new
PinotConfiguration(), tableConfig, SCHEMA, tableDataManager,
+ null));
+ assertTrue(e.getMessage().contains("must have a comparison column or a
time column configured"),
+ "Unexpected message: " + e.getMessage());
+ }
+
@Test
public void testCreateForManagerClassWithConsistentDeletes() {
UpsertConfig upsertConfig = new UpsertConfig(UpsertConfig.Mode.FULL);
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
index 120c89433ac..2354cd03486 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/utils/TableConfigUtilsTest.java
@@ -2594,6 +2594,7 @@ public class TableConfigUtilsTest {
@Test
public void testValidateUpsertConfig() {
UpsertConfig upsertConfig = new UpsertConfig(UpsertConfig.Mode.FULL);
+ upsertConfig.setComparisonColumn("myCol");
SegmentPartitionConfig segmentPartitionConfig =
new SegmentPartitionConfig(Map.of("myCol", new
ColumnPartitionConfig("murmur", 4)));
TableConfig validTableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
@@ -2621,6 +2622,30 @@ public class TableConfigUtilsTest {
+ "segment assignment. Configure segmentPartitionConfig in the
indexingConfig.");
}
+ // OFFLINE table with neither a comparison column nor a time column should
fail, same as realtime upsert
+ // (no implicit segment creation time fallback).
+ tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
+ .setUpsertConfig(new UpsertConfig(UpsertConfig.Mode.FULL))
+ .setSegmentPartitionConfig(segmentPartitionConfig)
+ .setRoutingConfig(STRICT_REPLICA_ROUTING_CONFIG)
+ .build();
+ try {
+ TableConfigUtils.validateUpsertAndDedupConfig(tableConfig, validSchema);
+ fail();
+ } catch (IllegalStateException e) {
+ assertEquals(e.getMessage(),
+ "Offline upsert table must have a comparison column or a time column
configured");
+ }
+
+ // OFFLINE table relying on the time column (no explicit comparison
column) should be allowed, same as realtime.
+ tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
+ .setTimeColumnName("myCol")
+ .setUpsertConfig(new UpsertConfig(UpsertConfig.Mode.FULL))
+ .setSegmentPartitionConfig(segmentPartitionConfig)
+ .setRoutingConfig(STRICT_REPLICA_ROUTING_CONFIG)
+ .build();
+ TableConfigUtils.validateUpsertAndDedupConfig(tableConfig, validSchema);
+
// OFFLINE table with partial upsert should fail
UpsertConfig partialUpsertConfig = new
UpsertConfig(UpsertConfig.Mode.PARTIAL);
tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(TABLE_NAME)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]