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]