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]

Reply via email to