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 968bf3cf618 [upsert] Support preload for OFFLINE table data managers 
(#19271)
968bf3cf618 is described below

commit 968bf3cf61819131e61b1860522576dfb71ea694
Author: Xiang Fu <[email protected]>
AuthorDate: Sun Aug 16 12:45:29 2026 -0700

    [upsert] Support preload for OFFLINE table data managers (#19271)
---
 .../apache/pinot/core/data/manager/BaseTableDataManager.java | 12 ++++++++++++
 .../core/data/manager/offline/OfflineTableDataManager.java   |  1 +
 .../core/data/manager/realtime/RealtimeTableDataManager.java | 12 ------------
 3 files changed, 13 insertions(+), 12 deletions(-)

diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
index df19704b855..4da4cf0d6ef 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/BaseTableDataManager.java
@@ -912,6 +912,18 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
     return Map.of();
   }
 
+  /// Handles upsert preload if the upsert preload is enabled.
+  protected void handleUpsertPreload(SegmentZKMetadata zkMetadata, 
IndexLoadingConfig indexLoadingConfig) {
+    if (_tableUpsertMetadataManager == null || 
!_tableUpsertMetadataManager.getContext().isPreloadEnabled()) {
+      return;
+    }
+    Integer partitionId = SegmentUtils.getSegmentPartitionId(zkMetadata, null);
+    Preconditions.checkState(partitionId != null,
+        "Failed to get partition id for segment: %s in upsert-enabled table: 
%s", zkMetadata.getSegmentName(),
+        _tableNameWithType);
+    
_tableUpsertMetadataManager.getOrCreatePartitionManager(partitionId).preloadSegments(indexLoadingConfig);
+  }
+
   protected void handleUpsert(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
     String segmentName = immutableSegment.getSegmentName();
     _logger.info("Adding immutable segment: {} with upsert enabled", 
segmentName);
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
index 66a8317f533..8717408446f 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/OfflineTableDataManager.java
@@ -73,6 +73,7 @@ public class OfflineTableDataManager extends 
BaseTableDataManager {
     SegmentZKMetadata zkMetadata = fetchZKMetadata(segmentName);
     IndexLoadingConfig indexLoadingConfig = fetchIndexLoadingConfig();
     indexLoadingConfig.setSegmentTier(zkMetadata.getTier());
+    handleUpsertPreload(zkMetadata, indexLoadingConfig);
     SegmentDataManager segmentDataManager = 
_segmentDataManagerMap.get(segmentName);
     if (segmentDataManager == null) {
       addNewOnlineSegment(zkMetadata, indexLoadingConfig);
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
index 9c85b0ad11c..934f14a9447 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManager.java
@@ -453,18 +453,6 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
     handleDedupPreload(zkMetadata, indexLoadingConfig);
   }
 
-  /// Handles upsert preload if the upsert preload is enabled.
-  private void handleUpsertPreload(SegmentZKMetadata zkMetadata, 
IndexLoadingConfig indexLoadingConfig) {
-    if (_tableUpsertMetadataManager == null || 
!_tableUpsertMetadataManager.getContext().isPreloadEnabled()) {
-      return;
-    }
-    Integer partitionId = SegmentUtils.getSegmentPartitionId(zkMetadata, null);
-    Preconditions.checkState(partitionId != null,
-        "Failed to get partition id for segment: %s in upsert-enabled table: 
%s", zkMetadata.getSegmentName(),
-        _tableNameWithType);
-    
_tableUpsertMetadataManager.getOrCreatePartitionManager(partitionId).preloadSegments(indexLoadingConfig);
-  }
-
   /// Handles dedup preload if the dedup preload is enabled.
   private void handleDedupPreload(SegmentZKMetadata zkMetadata, 
IndexLoadingConfig indexLoadingConfig) {
     if (_tableDedupMetadataManager == null || 
!_tableDedupMetadataManager.getContext().isPreloadEnabled()) {


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

Reply via email to