This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang 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 a3471b01c6b Give each upsert partition its own PartialUpsertHandler
(#19115)
a3471b01c6b is described below
commit a3471b01c6b646394071b6855e7f7cf76630eb95
Author: Kartik Khare <[email protected]>
AuthorDate: Sat Aug 1 04:05:50 2026 +0530
Give each upsert partition its own PartialUpsertHandler (#19115)
---
.../common/evaluator/InbuiltFunctionEvaluator.java | 5 ++++
.../recordtransformer/ExpressionTransformer.java | 5 ++++
.../upsert/BasePartitionUpsertMetadataManager.java | 6 +++-
.../upsert/BaseTableUpsertMetadataManager.java | 10 +++++--
.../segment/local/upsert/PartialUpsertHandler.java | 5 ++++
.../pinot/segment/local/upsert/UpsertContext.java | 32 ++++++++++++++--------
...rrentMapPartitionUpsertMetadataManagerTest.java | 6 ++--
7 files changed, 50 insertions(+), 19 deletions(-)
diff --git
a/pinot-common/src/main/java/org/apache/pinot/common/evaluator/InbuiltFunctionEvaluator.java
b/pinot-common/src/main/java/org/apache/pinot/common/evaluator/InbuiltFunctionEvaluator.java
index c30e1097a22..c72b1bdd4de 100644
---
a/pinot-common/src/main/java/org/apache/pinot/common/evaluator/InbuiltFunctionEvaluator.java
+++
b/pinot-common/src/main/java/org/apache/pinot/common/evaluator/InbuiltFunctionEvaluator.java
@@ -22,6 +22,7 @@ import com.google.common.base.Preconditions;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
+import javax.annotation.concurrent.NotThreadSafe;
import org.apache.commons.lang3.StringUtils;
import org.apache.pinot.common.function.FunctionInfo;
import org.apache.pinot.common.function.FunctionInvoker;
@@ -42,6 +43,10 @@ import org.apache.pinot.spi.function.FunctionEvaluator;
/// - FunctionNode - executes a function
/// - ColumnNode - fetches the value of the column from the input GenericRow
/// - ConstantNode - returns the literal value
+///
+/// NOTE: This class is not thread safe. Function nodes refill one reusable
argument array on every evaluation, so
+/// two threads evaluating the same instance can read each other's argument
values. Give each thread its own instance.
+@NotThreadSafe
public class InbuiltFunctionEvaluator implements FunctionEvaluator {
// Root of the execution tree
private final ExecutableNode _rootNode;
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/recordtransformer/ExpressionTransformer.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/recordtransformer/ExpressionTransformer.java
index ecc7928c2fd..1a3d457c4ed 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/recordtransformer/ExpressionTransformer.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/recordtransformer/ExpressionTransformer.java
@@ -28,6 +28,7 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import javax.annotation.Nullable;
+import javax.annotation.concurrent.NotThreadSafe;
import org.apache.pinot.common.evaluator.FunctionEvaluatorFactory;
import org.apache.pinot.common.utils.ThrottledLogger;
import org.apache.pinot.segment.local.utils.SchemaUtils;
@@ -47,7 +48,11 @@ import org.slf4j.LoggerFactory;
///
/// NOTE: should put this before the [DataTypeTransformer]. After this,
transformed column can be treated as
/// regular column for other record transformers.
+///
+/// NOTE: This class is not thread safe. It holds [FunctionEvaluator]
instances that keep reusable evaluation
+/// state, so each thread needs its own instance.
/// TODO: Merge this and CustomFunctionEnricher
+@NotThreadSafe
public class ExpressionTransformer implements RecordTransformer {
private static final Logger LOGGER =
LoggerFactory.getLogger(ExpressionTransformer.class);
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 02771620d3b..4c370cfb2d3 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
@@ -34,6 +34,7 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.Supplier;
import javax.annotation.Nullable;
import javax.annotation.concurrent.ThreadSafe;
import org.apache.helix.HelixManager;
@@ -155,7 +156,10 @@ public abstract class BasePartitionUpsertMetadataManager
implements PartitionUps
_comparisonColumns = context.getComparisonColumns();
_deleteRecordColumn = context.getDeleteRecordColumn();
_hashFunction = context.getHashFunction();
- _partialUpsertHandler = context.getPartialUpsertHandler();
+ // Build a handler owned by this partition. PartialUpsertHandler is not
thread safe, and merges for a partition
+ // run on one consumer thread at a time, same as the _reusePreviousRow
scratch state used alongside it.
+ Supplier<PartialUpsertHandler> partialUpsertHandlerSupplier =
context.getPartialUpsertHandlerSupplier();
+ _partialUpsertHandler = partialUpsertHandlerSupplier != null ?
partialUpsertHandlerSupplier.get() : null;
_enableSnapshot = context.isSnapshotEnabled();
_isPreloading = context.isPreloadEnabled();
_metadataTTL = context.getMetadataTTL();
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 b4d6ff162ec..7c801e4634a 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
@@ -20,6 +20,7 @@ package org.apache.pinot.segment.local.upsert;
import com.google.common.base.Preconditions;
import java.util.List;
+import java.util.function.Supplier;
import javax.annotation.Nullable;
import javax.annotation.concurrent.ThreadSafe;
import org.apache.commons.collections4.CollectionUtils;
@@ -68,9 +69,12 @@ public abstract class BaseTableUpsertMetadataManager
implements TableUpsertMetad
}
}
- PartialUpsertHandler partialUpsertHandler = null;
+ // PartialUpsertHandler is not thread safe, so hand each partition a
factory rather than one shared instance.
+ Supplier<PartialUpsertHandler> partialUpsertHandlerSupplier = null;
if (upsertConfig.getMode() == UpsertConfig.Mode.PARTIAL) {
- partialUpsertHandler = new PartialUpsertHandler(tableConfig, schema,
comparisonColumns, upsertConfig);
+ List<String> handlerComparisonColumns = comparisonColumns;
+ partialUpsertHandlerSupplier =
+ () -> new PartialUpsertHandler(tableConfig, schema,
handlerComparisonColumns, upsertConfig);
}
boolean enableSnapshot = upsertConfig.getSnapshot()
@@ -128,7 +132,7 @@ public abstract class BaseTableUpsertMetadataManager
implements TableUpsertMetad
.setPrimaryKeyColumns(primaryKeyColumns)
.setHashFunction(upsertConfig.getHashFunction())
.setComparisonColumns(comparisonColumns)
- .setPartialUpsertHandler(partialUpsertHandler)
+ .setPartialUpsertHandlerSupplier(partialUpsertHandlerSupplier)
.setDeleteRecordColumn(upsertConfig.getDeleteRecordColumn())
.setDropOutOfOrderRecord(upsertConfig.isDropOutOfOrderRecord())
.setOutOfOrderRecordColumn(upsertConfig.getOutOfOrderRecordColumn())
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
index 5a24408f890..d385706a25e 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/PartialUpsertHandler.java
@@ -22,6 +22,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import javax.annotation.Nullable;
+import javax.annotation.concurrent.NotThreadSafe;
import org.apache.pinot.segment.local.recordtransformer.RecordTransformerUtils;
import org.apache.pinot.segment.local.segment.readers.LazyRow;
import org.apache.pinot.segment.local.upsert.merger.PartialUpsertMerger;
@@ -42,6 +43,10 @@ import
org.apache.pinot.spi.recordtransformer.RecordTransformer;
///
/// It is also possible to define a custom logic for merging rows by
implementing [PartialUpsertMerger].
/// If a merger for row is defined then it takes precedence and ignores column
mergers.
+///
+/// NOTE: This class is not thread safe. The post partial upsert transformers
keep reusable evaluation state, so a
+/// handler belongs to exactly one partition and must be used by one thread at
a time.
+@NotThreadSafe
public class PartialUpsertHandler {
private final List<String> _primaryKeyColumns;
private final List<String> _comparisonColumns;
diff --git
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertContext.java
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertContext.java
index 25a90570948..625da707f94 100644
---
a/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertContext.java
+++
b/pinot-segment-local/src/main/java/org/apache/pinot/segment/local/upsert/UpsertContext.java
@@ -22,6 +22,7 @@ import com.google.common.base.Preconditions;
import java.io.File;
import java.util.List;
import java.util.Map;
+import java.util.function.Supplier;
import javax.annotation.Nullable;
import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.lang3.builder.ToStringBuilder;
@@ -41,8 +42,10 @@ public class UpsertContext {
private final List<String> _primaryKeyColumns;
private final HashFunction _hashFunction;
private final List<String> _comparisonColumns;
+ // Builds a handler for one partition. PartialUpsertHandler is not thread
safe, so the context hands out a factory
+ // instead of a shared instance. Null for full upsert.
@Nullable
- private final PartialUpsertHandler _partialUpsertHandler;
+ private final Supplier<PartialUpsertHandler> _partialUpsertHandlerSupplier;
@Nullable
private final String _deleteRecordColumn;
private final boolean _dropOutOfOrderRecord;
@@ -68,7 +71,8 @@ public class UpsertContext {
private final TableDataManager _tableDataManager;
private final File _tableIndexDir;
private UpsertContext(TableConfig tableConfig, Schema schema, List<String>
primaryKeyColumns,
- HashFunction hashFunction, List<String> comparisonColumns, @Nullable
PartialUpsertHandler partialUpsertHandler,
+ HashFunction hashFunction, List<String> comparisonColumns,
+ @Nullable Supplier<PartialUpsertHandler> partialUpsertHandlerSupplier,
@Nullable String deleteRecordColumn, boolean dropOutOfOrderRecord,
@Nullable String outOfOrderRecordColumn,
boolean enableSnapshot, boolean enablePreload, double metadataTTL,
double deletedKeysTTL,
boolean enableDeletedKeysCompactionConsistency,
UpsertConfig.ConsistencyMode consistencyMode,
@@ -80,7 +84,7 @@ public class UpsertContext {
_primaryKeyColumns = primaryKeyColumns;
_hashFunction = hashFunction;
_comparisonColumns = comparisonColumns;
- _partialUpsertHandler = partialUpsertHandler;
+ _partialUpsertHandlerSupplier = partialUpsertHandlerSupplier;
_deleteRecordColumn = deleteRecordColumn;
_dropOutOfOrderRecord = dropOutOfOrderRecord;
_outOfOrderRecordColumn = outOfOrderRecordColumn;
@@ -118,13 +122,14 @@ public class UpsertContext {
return _comparisonColumns;
}
+ /// Returns a factory that builds a [PartialUpsertHandler] for one
partition, or null for full upsert.
@Nullable
- public PartialUpsertHandler getPartialUpsertHandler() {
- return _partialUpsertHandler;
+ public Supplier<PartialUpsertHandler> getPartialUpsertHandlerSupplier() {
+ return _partialUpsertHandlerSupplier;
}
public UpsertConfig.Mode getUpsertMode() {
- return _partialUpsertHandler == null ? UpsertConfig.Mode.FULL :
UpsertConfig.Mode.PARTIAL;
+ return _partialUpsertHandlerSupplier == null ? UpsertConfig.Mode.FULL :
UpsertConfig.Mode.PARTIAL;
}
@Nullable
@@ -197,7 +202,7 @@ public class UpsertContext {
/// - Partial upsert is enabled (records need to be merged with previous
values)
/// - dropOutOfOrderRecord is enabled with NONE consistency mode (records
may have been dropped)
public boolean isTableTypeInconsistentDuringConsumption() {
- return _dropOutOfOrderRecord || _outOfOrderRecordColumn != null ||
_partialUpsertHandler != null;
+ return _dropOutOfOrderRecord || _outOfOrderRecordColumn != null ||
_partialUpsertHandlerSupplier != null;
}
@Override
@@ -230,7 +235,7 @@ public class UpsertContext {
private List<String> _primaryKeyColumns;
private HashFunction _hashFunction = HashFunction.NONE;
private List<String> _comparisonColumns;
- private PartialUpsertHandler _partialUpsertHandler;
+ private Supplier<PartialUpsertHandler> _partialUpsertHandlerSupplier;
private String _deleteRecordColumn;
private boolean _dropOutOfOrderRecord;
@Nullable
@@ -274,8 +279,10 @@ public class UpsertContext {
return this;
}
- public Builder setPartialUpsertHandler(PartialUpsertHandler
partialUpsertHandler) {
- _partialUpsertHandler = partialUpsertHandler;
+ /// Sets the factory used to build one [PartialUpsertHandler] per
partition. Null for full upsert.
+ public Builder setPartialUpsertHandlerSupplier(
+ @Nullable Supplier<PartialUpsertHandler> partialUpsertHandlerSupplier)
{
+ _partialUpsertHandlerSupplier = partialUpsertHandlerSupplier;
return this;
}
@@ -384,8 +391,9 @@ public class UpsertContext {
}
}
return new UpsertContext(_tableConfig, _schema, _primaryKeyColumns,
_hashFunction, _comparisonColumns,
- _partialUpsertHandler, _deleteRecordColumn, _dropOutOfOrderRecord,
_outOfOrderRecordColumn, _enableSnapshot,
- _enablePreload, _metadataTTL, _deletedKeysTTL,
_enableDeletedKeysCompactionConsistency, _consistencyMode,
+ _partialUpsertHandlerSupplier, _deleteRecordColumn,
_dropOutOfOrderRecord, _outOfOrderRecordColumn,
+ _enableSnapshot, _enablePreload, _metadataTTL, _deletedKeysTTL,
+ _enableDeletedKeysCompactionConsistency, _consistencyMode,
_upsertViewRefreshIntervalMs, _newSegmentTrackingTimeMs,
_metadataManagerConfigs,
_allowPartialUpsertConsumptionDuringCommit, _tableDataManager,
_tableIndexDir);
}
diff --git
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
index 198acb04bc8..a3dcfb12821 100644
---
a/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
+++
b/pinot-segment-local/src/test/java/org/apache/pinot/segment/local/upsert/ConcurrentMapPartitionUpsertMetadataManagerTest.java
@@ -2022,7 +2022,7 @@ public class
ConcurrentMapPartitionUpsertMetadataManagerTest {
// Test partial upserts with old and new segments having same number of
docs
// This test verifies that when all keys are present, no reversion occurs
PartialUpsertHandler mockPartialUpsertHandler =
mock(PartialUpsertHandler.class);
- UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandler(mockPartialUpsertHandler)
+ UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandlerSupplier(() -> mockPartialUpsertHandler)
.setConsistencyMode(UpsertConfig.ConsistencyMode.NONE).build();
ConcurrentMapPartitionUpsertMetadataManager upsertMetadataManager =
@@ -2085,7 +2085,7 @@ public class
ConcurrentMapPartitionUpsertMetadataManagerTest {
// Test partial upserts with consuming (mutable) segment being sealed -
revert should be triggered
// Note: Revert logic only applies when sealing a consuming segment, not
for immutable segment replacement
PartialUpsertHandler mockPartialUpsertHandler =
mock(PartialUpsertHandler.class);
- UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandler(mockPartialUpsertHandler)
+ UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandlerSupplier(() -> mockPartialUpsertHandler)
.setConsistencyMode(UpsertConfig.ConsistencyMode.NONE).build();
ConcurrentMapPartitionUpsertMetadataManager upsertMetadataManager =
@@ -2133,7 +2133,7 @@ public class
ConcurrentMapPartitionUpsertMetadataManagerTest {
public void testPartialUpsertOldSegmentLesserDocs() throws IOException {
// Test partial upserts with old segment having fewer docs than new segment
PartialUpsertHandler mockPartialUpsertHandler =
mock(PartialUpsertHandler.class);
- UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandler(mockPartialUpsertHandler)
+ UpsertContext upsertContext =
_contextBuilder.setPartialUpsertHandlerSupplier(() -> mockPartialUpsertHandler)
.setConsistencyMode(UpsertConfig.ConsistencyMode.NONE).build();
ConcurrentMapPartitionUpsertMetadataManager upsertMetadataManager =
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]