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 d4365386d8b Fix realtime table deletion races with segment creation 
(#19468)
d4365386d8b is described below

commit d4365386d8ba3605ec9b2fddd688a23c8ece80e5
Author: Xiang Fu <[email protected]>
AuthorDate: Thu Sep 24 22:37:24 2026 -0700

    Fix realtime table deletion races with segment creation (#19468)
    
    ## What this fixes
    
    Deleting a realtime table while a CONSUMING or ONLINE segment transition is 
in flight can leave a table manager or segment behind. Same-name recreation can 
also reuse a lease extender that the previous manager subsequently shuts down.
    
    This change closes the local deletion races together:
    
    - Record the newest deletion timestamp even when the server has not created 
a table manager yet. Late transitions must observe a missing config or a 
genuinely newer table incarnation; failed manager startup preserves this guard.
    - Coordinate consuming partition-manager creation and segment 
registration/start with shutdown. A segment that finishes construction after 
shutdown is destroyed without starting consumption or decrementing gauges for a 
segment that was never registered. A segment admitted before shutdown finishes 
starting before shutdown drains it.
    - Re-check shutdown before mutating the segment data directory. A stale 
consuming add of a deleted table aborts before deleting the directory, and a 
stale ONLINE download aborts before replacing it; a recreated same-name table's 
operations on a colliding segment name serialize behind the stale operation on 
the shared per-segment lock and find their files intact.
    - Admit immutable segment additions before metadata updates and 
registration. Shutdown rejects and destroys late-loaded segments, and waits for 
already admitted additions to finish before stopping partition managers and 
draining segments. This covers ONLINE loads, reloads, consuming-to-immutable 
replacement, upsert's two-stage replacement, dedup, and dimension lookup 
publication.
    - Prevent ONLINE preloading from creating upsert or dedup partition 
managers after shutdown has closed admission.
    - Hold a distinct per-table lifecycle lock through the old manager's full 
shutdown before permitting a replacement to initialize its lease extender. 
Locks use reclaimable weak values; hash-colliding table names remain 
independent.
    - Keep stale consuming segments from acting on a recreated same-name table 
after shutdown. A consuming segment whose constructor failed no longer posts 
(or keeps retrying) `segmentStoppedConsuming` once its table manager is shut 
down, so it cannot ask the controller to mark the recreated table's replica 
OFFLINE. The CONSUMING -> ONLINE transition holds a reference to the consuming 
segment so shutdown cannot destroy the mutable segment underneath the Helix 
thread's catch-up or build; s [...]
    
    Existing-manager lookups avoid the lifecycle lock. Downloads, construction, 
preloading, immutable-segment lifecycle callbacks, and blocking cleanup stay 
outside the segment-map admission monitor. Immutable additions can proceed 
concurrently; shutdown waits for the admitted operations rather than holding 
that monitor through their work. Segment-add callbacks remain outside the 
instance lifecycle lock.
    
    ## How to reproduce
    
    1. **Late first transition:** deliver a table-deletion message to a server 
with no local manager while the old table config is still present, then deliver 
an OFFLINE -> CONSUMING transition. Previously deletion recorded no timestamp, 
allowing the old config to recreate a manager.
    2. **Construction loses to deletion:** pause a consuming transition after 
construction but before registration. Delete the table, then resume. Previously 
it registered and started a consumer on the detached, shut-down manager.
    3. **Shutdown overtakes startup:** pause consumption startup after 
registration and delete the table concurrently. Previously shutdown could drain 
the segment before its consumer thread started.
    4. **Same-name recreation:** pause old-manager shutdown and concurrently 
initialize the replacement. Previously the replacement could reuse the lease 
extender that old shutdown subsequently removed.
    5. **Late ONLINE load:** pause ONLINE loading before immutable admission, 
delete the table, then resume. Previously the shutdown check threw without 
destroying the loaded segment.
    6. **ONLINE registration or upsert replacement:** pause immutable 
registration, or pause upsert replacement between its metadata and registration 
steps, then delete the table. Previously shutdown could drain before 
publication completed, leaving an unowned segment or racing the intermediate 
dual-segment manager.
    7. **Late ONLINE preload:** pause configuration loading before upsert/dedup 
preloading, delete the table, then resume. Previously preloading could create a 
partition manager after shutdown had stopped and closed the table's metadata 
managers.
    8. **Stale segment-directory mutation:** pause a consuming transition after 
admission but before the segment-directory cleanup, delete the table, then 
resume. Previously the stale add deleted `_indexDir/<segmentName>` even though 
a recreated same-name table (minute-precision LLC names restarting at partition 
sequence 0) could already own that directory. The ONLINE download path had the 
same window at `moveSegment`, which replaces the directory with the stale 
table's data.
    9. **Stale initialization failure:** let a consuming segment's constructor 
fail, delete the table, and recreate it with the same name so the LLC segment 
name collides. Previously the 30-second follow-up task found no manager 
registered on the old table and asked the controller to mark the recreated 
table's replica OFFLINE; if the first attempt was rejected it kept retrying 
until one was accepted.
    10. **Deletion during CONSUMING -> ONLINE:** pause `goOnlineFromConsuming` 
and delete the table. Previously shutdown released the map's only reference and 
destroyed the mutable segment while the transition was still reading it.
    11. **Stale COMMITTING download wait:** pause the pauseless download wait 
for a COMMITTING segment and delete the table. Previously the wait kept polling 
until the download timeout while holding the shared per-segment lock, blocking 
a recreated same-name table's operations on that segment.
    
    ## Validation
    
    - Rebased onto master at 257a0e7df2 (#19571 cached `IndexLoadingConfig`, 
#19444, #19629). `doAddConsumingSegment` now resolves its config through 
`getCachedIndexLoadingConfig()`, and the lifecycle test harness routes that 
accessor through `fetchIndexLoadingConfig()` so the ONLINE preload regression 
still pauses before preload.
    - 197 focused tests passed with no failures or skips: 
`HelixInstanceDataManagerLifecycleTest`, `HelixInstanceDataManagerTest`, 
`RealtimeTableDataManagerTest`, `RealtimeSegmentDataManagerTest`, 
`BaseTableDataManagerTest`, `BaseTableDataManagerNeedRefreshTest`, and 
`DimensionTableDataManagerTest`.
    - Every lifecycle regression fails with its guard removed and passes with 
it. This was re-verified after the rebase for the preload guards 
(`testOnlinePreloadDoesNotCreatePartitionAfterShutdown`) and for the four 
regressions added afterwards 
(`testInitializationErrorStopMsgSkippedAfterTableShutdown`, 
`testStopConsumedMsgRetryStopsAfterTableShutdown`, 
`testShutdownDoesNotDestroyConsumingSegmentDuringOnlineTransition`, 
`testCommittingSegmentDownloadAbortsAfterShutdown`).
    - Spotless, Checkstyle, license formatting/checks, and a 
deprecation-enabled compile of `pinot-core` and `pinot-server` (main and test) 
passed with no warnings on added lines. `git diff --check` is clean.
    - Independent review: every scenario above was traced as closed on the 
rebased code, and the findings that survived adversarial verification were 
either fixed here (items 9-11) or are the pre-existing follow-ups below.
    
    The deletion timestamp retains the existing controller-message/ZooKeeper 
creation-time contract. This change handles received deletion messages and 
local segment transitions; it does not change controller message delivery or 
introduce a new table-generation protocol.
    
    ## Known follow-ups (pre-existing, not changed here)
    
    - Upsert/dedup preload preprocesses `_indexDir/<segment>` in place on the 
preload executor with no per-segment lock and only an entry-time shutdown 
check, so a long reprocess of a deleted table's segment can overlap a 
same-minute name collision.
    - A delayed `TableDeletionMessage` can tear down a same-name recreated 
table's manager because `deleteTable` has no incarnation check, and the 
lock-free `getOrCreateTableDataManager` fast path can route a recreated table's 
first transition to the not-yet-deleted old manager. Both need the controller 
ordering or a manager-side config creation time and are left for a separate 
change.
---
 .../core/data/manager/BaseTableDataManager.java    |  52 +-
 .../core/data/manager/InstanceDataManager.java     |   4 +-
 .../manager/offline/DimensionTableDataManager.java |   4 +-
 .../manager/offline/OfflineTableDataManager.java   |   9 +-
 .../realtime/RealtimeSegmentDataManager.java       |  22 +-
 .../manager/realtime/RealtimeTableDataManager.java | 126 ++++-
 .../data/manager/BaseTableDataManagerTest.java     | 125 +++++
 .../realtime/RealtimeSegmentDataManagerTest.java   |  56 +-
 .../realtime/RealtimeTableDataManagerTest.java     | 610 +++++++++++++++++++++
 .../starter/helix/HelixInstanceDataManager.java    |  82 ++-
 .../HelixInstanceDataManagerLifecycleTest.java     | 380 +++++++++++++
 11 files changed, 1403 insertions(+), 67 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 da5fa943ec3..ec088917805 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
@@ -38,6 +38,7 @@ import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
+import java.util.concurrent.Phaser;
 import java.util.concurrent.Semaphore;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
@@ -134,6 +135,10 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
   protected static final Logger LOGGER = 
LoggerFactory.getLogger(BaseTableDataManager.class);
 
   protected final ConcurrentHashMap<String, SegmentDataManager> 
_segmentDataManagerMap = new ConcurrentHashMap<>();
+  /// Admission barrier for immutable segment additions. One party belongs to 
shutdown. Additions register under the
+  /// segment map monitor and finish without holding it, so metadata updates, 
lifecycle hooks and replacement cleanup
+  /// complete before shutdown stops the partition managers and drains the map.
+  private final Phaser _immutableSegmentAdds = new Phaser(1);
   protected final ServerMetrics _serverMetrics = ServerMetrics.get();
   protected TableUpsertMetadataManager _tableUpsertMetadataManager;
 
@@ -307,7 +312,12 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
       return;
     }
     _logger.info("Shutting down table data manager");
-    _shutDown = true;
+    // Close admission atomically with immutable additions and consuming 
segment startup. Downloads, construction
+    // and blocking cleanup must remain outside this monitor.
+    synchronized (_segmentDataManagerMap) {
+      _shutDown = true;
+    }
+    _immutableSegmentAdds.arriveAndAwaitAdvance();
     doShutdown();
     _logger.info("Shut down table data manager");
   }
@@ -370,9 +380,30 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
   /// @param immutableSegment Immutable segment to add
   @Override
   public void addSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
+    boolean admitted;
+    synchronized (_segmentDataManagerMap) {
+      admitted = !_shutDown;
+      if (admitted) {
+        _immutableSegmentAdds.register();
+      }
+    }
+    if (!admitted) {
+      // Loading can complete after shutdown. This segment was never 
registered or counted in table gauges.
+      immutableSegment.destroy();
+      throw new IllegalStateException("Table data manager is already shut 
down, cannot add segment: "
+          + immutableSegment.getSegmentName() + " to table: " + 
_tableNameWithType);
+    }
+    try {
+      doAddSegment(immutableSegment, zkMetadata);
+    } finally {
+      _immutableSegmentAdds.arriveAndDeregister();
+    }
+  }
+
+  /// Adds an admitted immutable segment, including its metadata and 
replacement cleanup. Shutdown waits for this
+  /// operation to finish before stopping partition managers and draining 
registered segments.
+  protected void doAddSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
     String segmentName = immutableSegment.getSegmentName();
-    Preconditions.checkState(!_shutDown, "Table data manager is already shut 
down, cannot add segment: %s to table: %s",
-        segmentName, _tableNameWithType);
     _logger.info("Adding immutable segment: {}", segmentName);
     _serverMetrics.addValueToTableGauge(_tableNameWithType, 
ServerGauge.DOCUMENT_COUNT,
         immutableSegment.getSegmentMetadata().getTotalDocs());
@@ -941,7 +972,13 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
     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);
+    PartitionUpsertMetadataManager partitionManager;
+    synchronized (_segmentDataManagerMap) {
+      Preconditions.checkState(!_shutDown, "Table data manager is already shut 
down, cannot preload table: %s",
+          _tableNameWithType);
+      partitionManager = 
_tableUpsertMetadataManager.getOrCreatePartitionManager(partitionId);
+    }
+    partitionManager.preloadSegments(indexLoadingConfig);
   }
 
   protected void handleUpsert(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
@@ -1488,6 +1525,13 @@ public abstract class BaseTableDataManager implements 
TableDataManager {
   protected File moveSegment(String segmentName, File untarredSegmentDir)
       throws IOException {
     File indexDir = getSegmentDataDir(segmentName);
+    // Replacing the segment data directory is destructive, and a deleted 
table shares the directory with a same-name
+    // recreated table. Abort stale operations instead of mutating the 
directory after shutdown. Callers hold the
+    // per-segment lock from the instance-wide SegmentLocks shared across 
table data managers, so the recreated
+    // table's lock-holding operations on the same segment cannot interleave 
with this method.
+    Preconditions.checkState(!_shutDown,
+        "Table data manager is already shut down, cannot replace data 
directory of segment: %s of table: %s",
+        segmentName, _tableNameWithType);
     try {
       FileUtils.deleteDirectory(indexDir);
       FileUtils.moveDirectory(untarredSegmentDir, indexDir);
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/InstanceDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/InstanceDataManager.java
index e71dfbff0fa..7e8647784bd 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/InstanceDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/InstanceDataManager.java
@@ -66,7 +66,9 @@ public interface InstanceDataManager {
   /// Should be called only once. After calling shut down, no other method 
should be called.
   void shutDown();
 
-  /// Delete a table.
+  /// Deletes a table: records `deletionTimeMs` as the table's newest deletion 
time, even when this instance never
+  /// created a data manager for it, then shuts down and removes the data 
manager if one exists. Until that record
+  /// expires, a same-name recreation is accepted only from a table config 
created after the recorded time.
   void deleteTable(String tableNameWithType, long deletionTimeMs)
       throws Exception;
 
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManager.java
index 858aa198814..19ee7943329 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/offline/DimensionTableDataManager.java
@@ -139,8 +139,8 @@ public class DimensionTableDataManager extends 
OfflineTableDataManager {
   }
 
   @Override
-  public void addSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
-    super.addSegment(immutableSegment, zkMetadata);
+  protected void doAddSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
+    super.doAddSegment(immutableSegment, zkMetadata);
     String segmentName = immutableSegment.getSegmentName();
     if (loadLookupTable()) {
       _logger.info("Successfully loaded lookup table after adding segment: 
{}", 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 74ed8c44e81..df98e48cfea 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
@@ -18,7 +18,6 @@
  */
 package org.apache.pinot.core.data.manager.offline;
 
-import com.google.common.base.Preconditions;
 import java.io.IOException;
 import javax.annotation.Nullable;
 import javax.annotation.concurrent.ThreadSafe;
@@ -82,16 +81,12 @@ public class OfflineTableDataManager extends 
BaseTableDataManager {
   }
 
   @Override
-  public void addSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
-    String segmentName = immutableSegment.getSegmentName();
-    Preconditions.checkState(!_shutDown,
-        "Table data manager is already shut down, cannot add segment: %s to 
table: %s",
-        segmentName, _tableNameWithType);
+  protected void doAddSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
     if (isUpsertEnabled()) {
       handleUpsert(immutableSegment, zkMetadata);
       return;
     }
-    super.addSegment(immutableSegment, zkMetadata);
+    super.doAddSegment(immutableSegment, zkMetadata);
   }
 
   @Override
diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
index 4c511c00b2d..1e239ce14da 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManager.java
@@ -1547,13 +1547,33 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
       if (response.getStatus() == 
SegmentCompletionProtocol.ControllerResponseStatus.PROCESSED) {
         break;
       }
+      // Nothing stops a segment that never started consuming, so re-check the 
table between retries: once it is
+      // shut down, a retry could be accepted on behalf of a recreated 
same-name table's segment.
+      if (isTableDataManagerShutDown()) {
+        _segmentLogger.info("Stop retrying segmentStoppedConsuming for 
segment: {}, table data manager is already "
+            + "shut down", _segmentNameStr);
+        break;
+      }
       Uninterruptibles.sleepUninterruptibly(10, TimeUnit.SECONDS);
       _segmentLogger.info("Retrying after response {}", 
response.toJsonString());
     } while (!_shouldStop);
   }
 
+  /// Whether the owning table data manager has been shut down (the table was 
deleted, or the server is stopping).
+  private boolean isTableDataManagerShutDown() {
+    return _realtimeTableDataManager != null && 
_realtimeTableDataManager.isShutDown();
+  }
+
   @VisibleForTesting
   void postStopConsumedMsgForInitializationError() {
+    if (isTableDataManagerShutDown()) {
+      // The table was shut down after initialization failed. Its Helix state 
no longer matters, and a recreated table
+      // with the same name may already own this segment name, so asking the 
controller to mark the segment OFFLINE
+      // could deregister the new table's consuming replica instead.
+      _segmentLogger.info("Skip segmentStoppedConsuming for segment: {}, table 
data manager is already shut down",
+          _segmentNameStr);
+      return;
+    }
     if (hasDifferentSegmentDataManagerRegistered()) {
       _segmentLogger.info(
           "Skip segmentStoppedConsuming for segment: {}, another segment data 
manager is already registered",
@@ -1816,7 +1836,7 @@ public class RealtimeSegmentDataManager extends 
SegmentDataManager {
   public void stop()
       throws InterruptedException {
     _shouldStop = true;
-    if (Thread.currentThread() != _consumerThread && 
_consumerThread.isAlive()) {
+    if (_consumerThread != null && Thread.currentThread() != _consumerThread 
&& _consumerThread.isAlive()) {
       _segmentLogger.info("Interrupting the consumer thread and waiting for it 
to join");
       long startTimeMs = System.currentTimeMillis();
       _consumerThread.interrupt();
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 6970d66d105..faad4713b2b 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
@@ -462,7 +462,13 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
     Preconditions.checkState(partitionId != null,
         "Failed to get partition id for segment: %s in dedup-enabled table: 
%s", zkMetadata.getSegmentName(),
         _tableNameWithType);
-    
_tableDedupMetadataManager.getOrCreatePartitionManager(partitionId).preloadSegments(indexLoadingConfig);
+    PartitionDedupMetadataManager partitionManager;
+    synchronized (_segmentDataManagerMap) {
+      Preconditions.checkState(!_shutDown, "Table data manager is already shut 
down, cannot preload table: %s",
+          _tableNameWithType);
+      partitionManager = 
_tableDedupMetadataManager.getOrCreatePartitionManager(partitionId);
+    }
+    partitionManager.preloadSegments(indexLoadingConfig);
   }
 
   protected void doAddOnlineSegment(String segmentName)
@@ -477,7 +483,17 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
       addNewOnlineSegment(zkMetadata, indexLoadingConfig);
     } else if (segmentDataManager instanceof RealtimeSegmentDataManager) {
       _logger.info("Changing segment: {} from CONSUMING to ONLINE", 
segmentName);
-      ((RealtimeSegmentDataManager) 
segmentDataManager).goOnlineFromConsuming(zkMetadata);
+      // Table shutdown can drain and release the consuming segment while this 
thread catches it up or builds it.
+      // Hold a reference for the transition so the mutable segment is 
destroyed by the last release rather than
+      // underneath the build; shutdown still offloads it without waiting for 
the transition to finish.
+      Preconditions.checkState(segmentDataManager.increaseReferenceCount(),
+          "Table data manager is already shut down, cannot change segment: %s 
from CONSUMING to ONLINE in table: %s",
+          segmentName, _tableNameWithType);
+      try {
+        ((RealtimeSegmentDataManager) 
segmentDataManager).goOnlineFromConsuming(zkMetadata);
+      } finally {
+        releaseSegment(segmentDataManager);
+      }
     } else if (zkMetadata.getStatus().isCompleted()) {
       // For pauseless ingestion, the segment is marked ONLINE before it's 
built and before the COMMIT_END_METADATA
       // call completes.
@@ -541,12 +557,45 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
       return;
     }
     IndexLoadingConfig indexLoadingConfig = getCachedIndexLoadingConfig();
-    handleSegmentPreload(zkMetadata, indexLoadingConfig);
-    SegmentDataManager segmentDataManager = 
_segmentDataManagerMap.get(segmentName);
-    if (segmentDataManager != null) {
-      _logger.warn("Segment: {} ({}) already exists, skipping adding it as 
CONSUMING segment", segmentName,
-          segmentDataManager instanceof RealtimeSegmentDataManager ? 
"CONSUMING" : "COMPLETED");
-      return;
+    LLCSegmentName llcSegmentName = new LLCSegmentName(segmentName);
+    int partitionGroupId = llcSegmentName.getPartitionGroupId();
+    PartitionUpsertMetadataManager partitionUpsertMetadataManager;
+    PartitionDedupMetadataManager partitionDedupMetadataManager;
+    synchronized (_segmentDataManagerMap) {
+      Preconditions.checkState(!_shutDown,
+          "Table data manager is already shut down, cannot add CONSUMING 
segment: %s to table: %s", segmentName,
+          _tableNameWithType);
+      // Partition owners must be visible before shutdown stops/closes them. 
Preload and construction can be slow,
+      // so run them outside this monitor and check admission again before 
publishing the segment.
+      partitionUpsertMetadataManager = _tableUpsertMetadataManager != null
+          ? 
_tableUpsertMetadataManager.getOrCreatePartitionManager(partitionGroupId)
+          : null;
+      partitionDedupMetadataManager = _tableDedupMetadataManager != null
+          ? 
_tableDedupMetadataManager.getOrCreatePartitionManager(partitionGroupId)
+          : null;
+    }
+    if (partitionUpsertMetadataManager != null && 
_tableUpsertMetadataManager.getContext().isPreloadEnabled()) {
+      partitionUpsertMetadataManager.preloadSegments(indexLoadingConfig);
+    }
+    if (partitionDedupMetadataManager != null && 
_tableDedupMetadataManager.getContext().isPreloadEnabled()) {
+      partitionDedupMetadataManager.preloadSegments(indexLoadingConfig);
+    }
+    beforeConsumingSegmentAdmissionRecheck(segmentName);
+    synchronized (_segmentDataManagerMap) {
+      // Shutdown does not wait for in-flight consuming adds, and after the 
table is deleted a same-name table can be
+      // recreated over the same data directory. The per-segment lock (from 
the instance-wide SegmentLocks shared
+      // across table data managers) keeps the recreated table's writers for 
this segment out until this add returns;
+      // this re-check additionally aborts a stale add before it deletes the 
previous incarnation's leftover directory
+      // or constructs a segment data manager that shutdown would only reject 
at publish time.
+      Preconditions.checkState(!_shutDown,
+          "Table data manager is already shut down, cannot add CONSUMING 
segment: %s to table: %s", segmentName,
+          _tableNameWithType);
+      SegmentDataManager segmentDataManager = 
_segmentDataManagerMap.get(segmentName);
+      if (segmentDataManager != null) {
+        _logger.warn("Segment: {} ({}) already exists, skipping adding it as 
CONSUMING segment", segmentName,
+            segmentDataManager instanceof RealtimeSegmentDataManager ? 
"CONSUMING" : "COMPLETED");
+        return;
+      }
     }
 
     _logger.info("Adding new CONSUMING segment: {}", segmentName);
@@ -565,29 +614,51 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
     setDefaultTimeValueIfInvalid(tableConfig, schema, zkMetadata);
 
     // Generates only one semaphore for every partition
-    LLCSegmentName llcSegmentName = new LLCSegmentName(segmentName);
-    int partitionGroupId = llcSegmentName.getPartitionGroupId();
     ConsumerCoordinator consumerCoordinator = 
getConsumerCoordinator(partitionGroupId);
 
     // Create the segment data manager and register it
-    PartitionUpsertMetadataManager partitionUpsertMetadataManager =
-        _tableUpsertMetadataManager != null ? 
_tableUpsertMetadataManager.getOrCreatePartitionManager(partitionGroupId)
-            : null;
-    PartitionDedupMetadataManager partitionDedupMetadataManager =
-        _tableDedupMetadataManager != null ? 
_tableDedupMetadataManager.getOrCreatePartitionManager(partitionGroupId)
-            : null;
     RealtimeSegmentDataManager realtimeSegmentDataManager =
         createRealtimeSegmentDataManager(zkMetadata, tableConfig, 
indexLoadingConfig, schema, llcSegmentName,
             consumerCoordinator, partitionUpsertMetadataManager, 
partitionDedupMetadataManager,
             _isTableReadyToConsumeData);
-    registerSegment(segmentName, realtimeSegmentDataManager, 
partitionUpsertMetadataManager);
-    if (partitionUpsertMetadataManager != null) {
-      partitionUpsertMetadataManager.trackNewlyAddedSegment(segmentName);
+    // A segment can finish construction after table shutdown has drained the 
registered segments. Publish and start
+    // it together so shutdown either owns its cleanup or rejects it before 
any consumer is started.
+    // Query threads never take this monitor: they read the segment map 
lock-free and only take per-segment reference
+    // counts and the upsert view locks. The upsert view locks acquired below 
(trackSegmentForUpsertView) are leaf
+    // locks that never call back into the table data manager, so holding this 
monitor while taking them cannot
+    // delay a query. Keep it that way: do not add callbacks from upsert view 
code into this class.
+    synchronized (_segmentDataManagerMap) {
+      if (!_shutDown) {
+        registerSegment(segmentName, realtimeSegmentDataManager, 
partitionUpsertMetadataManager);
+        if (partitionUpsertMetadataManager != null) {
+          partitionUpsertMetadataManager.trackNewlyAddedSegment(segmentName);
+        }
+        realtimeSegmentDataManager.startConsumption();
+        incrementSegmentCountGauge();
+        _logger.info("Added new CONSUMING segment: {}", segmentName);
+        return;
+      }
     }
-    realtimeSegmentDataManager.startConsumption();
+    // It was never published or counted in table gauges, so do not use 
releaseSegment/closeSegment here.
+    realtimeSegmentDataManager.destroy();
+    throw new IllegalStateException(
+        "Table data manager is already shut down, cannot add CONSUMING 
segment: " + segmentName + " to table: "
+            + _tableNameWithType);
+  }
+
+  /// Counts a published CONSUMING segment through the same deprecated 
`AbstractMetrics.addValueToTableGauge` API that
+  /// `BaseTableDataManager.closeSegment` uses for the matching decrement, so 
the SEGMENT_COUNT gauge stays balanced
+  /// until both sides migrate together.
+  @SuppressWarnings("deprecation")
+  private void incrementSegmentCountGauge() {
     _serverMetrics.addValueToTableGauge(_tableNameWithType, 
ServerGauge.SEGMENT_COUNT, 1);
+  }
 
-    _logger.info("Added new CONSUMING segment: {}", segmentName);
+  /// Invoked while adding a CONSUMING segment, after the initial admission 
and partition-manager preload, immediately
+  /// before the shutdown re-check that guards the segment data directory 
cleanup. The per-segment lock is held.
+  /// No-op in production; tests override it to pause an in-flight add inside 
the shutdown/recreation window.
+  @VisibleForTesting
+  void beforeConsumingSegmentAdmissionRecheck(String segmentName) {
   }
 
   @Override
@@ -606,6 +677,11 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
     long downloadTimeoutMs = 
getDownloadTimeoutMs(getCachedTableConfigAndSchema().getLeft());
     long deadlineMs = System.currentTimeMillis() + downloadTimeoutMs;
     while (System.currentTimeMillis() < deadlineMs) {
+      // A deleted table's transition can wait here for minutes while holding 
the shared per-segment lock. Abort so a
+      // recreated same-name table's operations on this segment are not 
blocked behind, or served by, a stale wait.
+      Preconditions.checkState(!_shutDown,
+          "Table data manager is already shut down, cannot download COMMITTING 
segment: %s of table: %s", segmentName,
+          _tableNameWithType);
       // ZK Metadata may change during segment download process; fetch it on 
every retry.
       zkMetadata = fetchZKMetadata(segmentName);
       if (zkMetadata.getStatus().isCompleted()) {
@@ -700,11 +776,7 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
   }
 
   @Override
-  public void addSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
-    String segmentName = immutableSegment.getSegmentName();
-    Preconditions.checkState(!_shutDown, "Table data manager is already shut 
down, cannot add segment: %s to table: %s",
-        segmentName, _tableNameWithType);
-
+  protected void doAddSegment(ImmutableSegment immutableSegment, @Nullable 
SegmentZKMetadata zkMetadata) {
     if (isUpsertEnabled()) {
       handleUpsert(immutableSegment, zkMetadata);
       return;
@@ -715,7 +787,7 @@ public class RealtimeTableDataManager extends 
BaseTableDataManager {
       handleDedup((ImmutableSegmentImpl) immutableSegment);
     }
 
-    super.addSegment(immutableSegment, zkMetadata);
+    super.doAddSegment(immutableSegment, zkMetadata);
   }
 
   private void handleDedup(ImmutableSegmentImpl immutableSegment) {
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
index c7faa46cd1a..2bb3267ab54 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/BaseTableDataManagerTest.java
@@ -30,14 +30,17 @@ import java.util.List;
 import java.util.Map;
 import java.util.UUID;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.FutureTask;
 import java.util.concurrent.TimeUnit;
 import org.apache.commons.io.FileUtils;
 import org.apache.helix.HelixManager;
 import org.apache.pinot.common.metadata.ZKMetadataProvider;
 import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
+import org.apache.pinot.common.metrics.ServerGauge;
 import org.apache.pinot.common.metrics.ServerMetrics;
 import org.apache.pinot.common.tier.TierFactory;
 import org.apache.pinot.common.utils.TarCompressionUtils;
@@ -89,6 +92,9 @@ import org.testng.annotations.BeforeMethod;
 import org.testng.annotations.Test;
 
 import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.clearInvocations;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
@@ -163,6 +169,100 @@ public class BaseTableDataManagerTest {
     PinotCrypterFactory.init(new PinotConfiguration(properties));
   }
 
+  @Test
+  @SuppressWarnings("deprecation")
+  public void testOnlineSegmentLoadedDuringShutdownIsDestroyed()
+      throws Exception {
+    BaseTableDataManager table = spy(createTableManager());
+    ServerMetrics metrics = table._serverMetrics;
+    clearInvocations(metrics);
+    ImmutableSegment segment = mock(ImmutableSegment.class);
+    when(segment.getSegmentName()).thenReturn(SEGMENT_NAME);
+    CountDownLatch loading = new CountDownLatch(1);
+    CountDownLatch finishLoading = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      loading.countDown();
+      assertTrue(finishLoading.await(10, TimeUnit.SECONDS));
+      table.addSegment(segment, null);
+      return null;
+    }).when(table).doAddOnlineSegment(SEGMENT_NAME);
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addOnlineSegment(SEGMENT_NAME);
+      return null;
+    });
+    Thread additionThread = new Thread(addition, "load-online-segment");
+    try {
+      additionThread.start();
+      assertTrue(loading.await(10, TimeUnit.SECONDS));
+      table.shutDown();
+      finishLoading.countDown();
+      ExecutionException failure = expectThrows(ExecutionException.class, () 
-> addition.get(10, TimeUnit.SECONDS));
+      assertTrue(failure.getCause() instanceof IllegalStateException);
+      assertEquals(table.getNumSegments(), 0);
+      verify(segment).destroy();
+      verify(metrics, never()).addValueToTableGauge(eq(OFFLINE_TABLE_NAME), 
eq(ServerGauge.SEGMENT_COUNT), anyLong());
+      verify(metrics, never()).addValueToTableGauge(eq(OFFLINE_TABLE_NAME), 
eq(ServerGauge.DOCUMENT_COUNT), anyLong());
+    } finally {
+      finishLoading.countDown();
+      additionThread.join(10000);
+    }
+  }
+
+  @Test
+  @SuppressWarnings("deprecation")
+  public void testShutdownWaitsForImmutableSegmentRegistration()
+      throws Exception {
+    BaseTableDataManager table = spy(createTableManager());
+    ServerMetrics metrics = table._serverMetrics;
+    clearInvocations(metrics);
+    ImmutableSegment segment = mock(ImmutableSegment.class);
+    when(segment.getSegmentName()).thenReturn(SEGMENT_NAME);
+    SegmentMetadata metadata = mock(SegmentMetadata.class);
+    when(metadata.getTotalDocs()).thenReturn(5);
+    when(segment.getSegmentMetadata()).thenReturn(metadata);
+    CountDownLatch registering = new CountDownLatch(1);
+    CountDownLatch finishRegistration = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      registering.countDown();
+      assertTrue(finishRegistration.await(10, TimeUnit.SECONDS));
+      return invocation.callRealMethod();
+    }).when(table).registerSegment(eq(SEGMENT_NAME), any());
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addSegment(segment, null);
+      return null;
+    });
+    FutureTask<Void> shutdown = new FutureTask<>(() -> {
+      table.shutDown();
+      return null;
+    });
+    Thread additionThread = new Thread(addition, "register-immutable-segment");
+    Thread shutdownThread = new Thread(shutdown, "shutdown-immutable-table");
+    try {
+      additionThread.start();
+      assertTrue(registering.await(10, TimeUnit.SECONDS));
+      shutdownThread.start();
+      TestUtils.waitForCondition(ignored -> shutdown.isDone() || 
(table.isShutDown()
+          && (shutdownThread.getState() == Thread.State.WAITING
+          || shutdownThread.getState() == Thread.State.TIMED_WAITING)), 10, 
10000,
+          "Shutdown did not wait for the admitted segment");
+      assertFalse(shutdown.isDone(), "Shutdown must wait for the admitted 
segment to finish registration");
+      verify(segment, never()).destroy();
+      finishRegistration.countDown();
+      addition.get(10, TimeUnit.SECONDS);
+      shutdown.get(10, TimeUnit.SECONDS);
+      assertEquals(table.getNumSegments(), 0);
+      verify(segment).destroy();
+      verify(metrics).addValueToTableGauge(OFFLINE_TABLE_NAME, 
ServerGauge.SEGMENT_COUNT, 1L);
+      verify(metrics).addValueToTableGauge(OFFLINE_TABLE_NAME, 
ServerGauge.SEGMENT_COUNT, -1L);
+      verify(metrics).addValueToTableGauge(OFFLINE_TABLE_NAME, 
ServerGauge.DOCUMENT_COUNT, 5L);
+      verify(metrics).addValueToTableGauge(OFFLINE_TABLE_NAME, 
ServerGauge.DOCUMENT_COUNT, -5L);
+    } finally {
+      finishRegistration.countDown();
+      additionThread.join(10000);
+      shutdownThread.join(10000);
+    }
+  }
+
   @Test
   public void testReloadSegmentNewData()
       throws Exception {
@@ -867,6 +967,31 @@ public class BaseTableDataManagerTest {
     }
   }
 
+  // Regression for the same-name table recreation race on the ONLINE path: 
after shutdown, a stale download must not
+  // replace the segment data directory, which a recreated same-name table may 
already own.
+  @Test
+  public void testMoveSegmentRejectedAfterShutdown()
+      throws IOException {
+    BaseTableDataManager tableDataManager = createTableManager();
+    File tempRootDir = 
tableDataManager.getTmpSegmentDataDir("test-move-after-shutdown");
+
+    File tempTar = new File(tempRootDir, SEGMENT_NAME + 
TarCompressionUtils.TAR_COMPRESSED_FILE_EXTENSION);
+    File tempInputDir = new File(tempRootDir, "input");
+    FileUtils.write(new File(tempInputDir, "tmp.txt"), "this is in segment 
dir", StandardCharsets.UTF_8);
+    TarCompressionUtils.createCompressedTarFile(tempInputDir, tempTar);
+    FileUtils.deleteQuietly(tempInputDir);
+
+    File segmentDir = tableDataManager.getSegmentDataDir(SEGMENT_NAME);
+    File marker = new File(segmentDir, "marker");
+    FileUtils.write(marker, "recreated owner's data", StandardCharsets.UTF_8);
+
+    tableDataManager.shutDown();
+    expectThrows(IllegalStateException.class,
+        () -> tableDataManager.untarAndMoveSegment(SEGMENT_NAME, tempTar, 
tempRootDir));
+    assertEquals(FileUtils.readFileToString(marker, StandardCharsets.UTF_8), 
"recreated owner's data",
+        "Stale download must not replace the segment data directory after 
shutdown");
+  }
+
   @Test
   public void 
testReplaceSegmentIfCrcMismatchWhenFlagDisabledSegmentCrcMismatchShouldDownload()
       throws Exception {
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
index fb75d2ef714..7aec91dc0d2 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeSegmentDataManagerTest.java
@@ -173,6 +173,49 @@ public class RealtimeSegmentDataManagerTest {
     }
   }
 
+  @Test
+  public void testInitializationErrorStopMsgSkippedAfterTableShutdown()
+      throws Exception {
+    try (FakeRealtimeSegmentDataManager segmentDataManager = 
createFakeSegmentManager()) {
+      RealtimeTableDataManager tableDataManager = 
segmentDataManager.getTableDataManager();
+      // Shutdown drained the old table's segment map, so nothing is 
registered under this segment name even though a
+      // recreated same-name table may already be consuming it.
+      
when(tableDataManager.getSegmentDataManager(SEGMENT_NAME_STR)).thenReturn(null);
+      when(tableDataManager.isShutDown()).thenReturn(true);
+
+      segmentDataManager.postStopConsumedMsgForInitializationError();
+
+      Assert.assertFalse(segmentDataManager._postConsumeStoppedCalled);
+    }
+  }
+
+  @Test(timeOut = 10_000)
+  public void testStopConsumedMsgRetryStopsAfterTableShutdown()
+      throws Exception {
+    try (FakeRealtimeSegmentDataManager segmentDataManager = 
createFakeSegmentManager()) {
+      RealtimeTableDataManager tableDataManager = 
segmentDataManager.getTableDataManager();
+      // The controller rejects the first attempt (e.g. the table's ideal 
state is already gone) and the table is shut
+      // down before the retry: the loop must give up instead of sleeping and 
retrying on behalf of a recreated table.
+      segmentDataManager._stopConsumedResponseStatus = 
SegmentCompletionProtocol.ControllerResponseStatus.FAILED;
+      when(tableDataManager.isShutDown()).thenReturn(true);
+
+      segmentDataManager.invokePostStopConsumedMsg("Consuming segment 
initialization error");
+
+      Assert.assertEquals(segmentDataManager._postConsumeStoppedAttempts, 1);
+    }
+  }
+
+  @Test
+  public void testDestroyBeforeConsumptionStarts()
+      throws Exception {
+    FakeRealtimeSegmentDataManager segmentDataManager = 
createFakeSegmentManager();
+    // Never started: the real stop() must tolerate the missing consumer 
thread (the fake normally overrides it), and
+    // destroy() must still close the stream consumer.
+    segmentDataManager.invokeRealStop();
+    segmentDataManager.destroy();
+    Assert.assertTrue(segmentDataManager.isStreamConsumerClosed());
+  }
+
   @Test
   public void testPostStopConsumedMsgDoesNotCheckRegisteredSegmentManager()
       throws Exception {
@@ -1469,6 +1512,9 @@ public class RealtimeSegmentDataManagerTest {
     private boolean _notifySegmentBuildFailedWithDeterministicErrorCalled = 
false;
     public boolean _throwExceptionFromConsume = false;
     public boolean _postConsumeStoppedCalled = false;
+    public int _postConsumeStoppedAttempts = 0;
+    public SegmentCompletionProtocol.ControllerResponseStatus 
_stopConsumedResponseStatus =
+        SegmentCompletionProtocol.ControllerResponseStatus.PROCESSED;
     public Map<Integer, ConsumerCoordinator> _consumerCoordinatorMap;
     public boolean _stubConsumeLoop = true;
     public RealtimeTableDataManager _tableDataManager;
@@ -1606,12 +1652,18 @@ public class RealtimeSegmentDataManagerTest {
       super.postStopConsumedMsg(reason);
     }
 
+    /// Runs the production stop() instead of this fake's override.
+    public void invokeRealStop()
+        throws InterruptedException {
+      super.stop();
+    }
+
     @Override
     SegmentCompletionProtocol.Response 
postSegmentStoppedConsuming(ConsumptionStopIndicator indicator) {
       _postConsumeStoppedCalled = true;
+      _postConsumeStoppedAttempts++;
       return new SegmentCompletionProtocol.Response(
-          new SegmentCompletionProtocol.Response.Params().withStatus(
-              SegmentCompletionProtocol.ControllerResponseStatus.PROCESSED));
+          new 
SegmentCompletionProtocol.Response.Params().withStatus(_stopConsumedResponseStatus));
     }
 
     @Override
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManagerTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManagerTest.java
index f2e3955d147..e39b0151dee 100644
--- 
a/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManagerTest.java
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/data/manager/realtime/RealtimeTableDataManagerTest.java
@@ -21,33 +21,643 @@ package org.apache.pinot.core.data.manager.realtime;
 import com.google.common.cache.Cache;
 import com.google.common.cache.CacheBuilder;
 import com.google.common.cache.RemovalNotification;
+import java.io.File;
+import java.nio.file.Files;
 import java.time.Duration;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.FutureTask;
 import java.util.concurrent.Semaphore;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.BooleanSupplier;
+import org.apache.commons.io.FileUtils;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.commons.lang3.tuple.Pair;
 import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
+import org.apache.pinot.common.metrics.ServerGauge;
+import org.apache.pinot.common.metrics.ServerMetrics;
+import org.apache.pinot.common.utils.LLCSegmentName;
 import org.apache.pinot.core.realtime.impl.fakestream.FakeStreamConfigUtils;
+import org.apache.pinot.segment.local.dedup.DedupContext;
+import org.apache.pinot.segment.local.dedup.PartitionDedupMetadataManager;
+import org.apache.pinot.segment.local.dedup.TableDedupMetadataManager;
+import org.apache.pinot.segment.local.segment.index.loader.IndexLoadingConfig;
+import org.apache.pinot.segment.local.upsert.PartitionUpsertMetadataManager;
+import org.apache.pinot.segment.local.upsert.TableUpsertMetadataManager;
+import org.apache.pinot.segment.local.upsert.UpsertContext;
+import org.apache.pinot.segment.local.utils.SegmentLocks;
+import org.apache.pinot.segment.spi.ImmutableSegment;
+import org.apache.pinot.segment.spi.MutableSegment;
+import org.apache.pinot.segment.spi.SegmentMetadata;
 import org.apache.pinot.spi.config.table.TableConfig;
 import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.config.table.UpsertConfig.ConsistencyMode;
 import org.apache.pinot.spi.data.DateTimeFieldSpec;
 import org.apache.pinot.spi.data.FieldSpec;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
 import org.apache.pinot.spi.data.Schema;
 import org.apache.pinot.spi.stream.StreamConfig;
 import org.apache.pinot.spi.stream.StreamConsumerFactoryProvider;
 import org.apache.pinot.spi.stream.StreamMetadataProvider;
+import org.apache.pinot.spi.utils.CommonConstants;
 import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
 import org.apache.pinot.util.TestUtils;
 import org.joda.time.DateTimeZone;
 import org.joda.time.format.DateTimeFormat;
 import org.mockito.Mockito;
+import org.slf4j.LoggerFactory;
 import org.testng.Assert;
+import org.testng.annotations.DataProvider;
 import org.testng.annotations.Test;
 
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doCallRealMethod;
+import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
 import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertTrue;
+import static org.testng.Assert.expectThrows;
 
 
 public class RealtimeTableDataManagerTest {
+  @Test
+  public void testStopBeforeConsumptionStarts()
+      throws Exception {
+    RealtimeSegmentDataManager segment = 
mock(RealtimeSegmentDataManager.class);
+    doCallRealMethod().when(segment).stop();
+    segment.stop();
+  }
+
+  @DataProvider
+  public Object[][] upsertEnabled() {
+    return new Object[][]{{false}, {true}};
+  }
+
+  @Test(dataProvider = "upsertEnabled")
+  // Verify the existing gauge API used by consuming admission and segment 
cleanup.
+  @SuppressWarnings("deprecation")
+  public void testConsumingSegmentConstructedDuringShutdownIsDestroyed(boolean 
upsertEnabled)
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    ServerMetrics metrics = ServerMetrics.get();
+    clearInvocations(metrics);
+    File indexDir = Files.createTempDirectory("consuming-shutdown").toFile();
+    RealtimeSegmentDataManager segment = 
mock(RealtimeSegmentDataManager.class);
+    CountDownLatch constructed = new CountDownLatch(1);
+    CountDownLatch finishConstruction = new CountDownLatch(1);
+    LifecycleTableDataManager table = new LifecycleTableDataManager(indexDir, 
segment, upsertEnabled, () -> {
+      constructed.countDown();
+      await(finishConstruction);
+    });
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    Thread additionThread = new Thread(addition, 
"construct-consuming-segment");
+    try {
+      additionThread.start();
+      await(constructed, addition);
+      table.shutDown();
+      finishConstruction.countDown();
+      ExecutionException failure = expectThrows(ExecutionException.class, () 
-> addition.get(10, TimeUnit.SECONDS));
+      assertTrue(failure.getCause() instanceof IllegalStateException);
+      assertEquals(table.getNumSegments(), 0);
+      verify(segment).destroy();
+      verify(segment, never()).startConsumption();
+      verify(metrics, never()).addValueToTableGauge(eq(table.getTableName()), 
eq(ServerGauge.SEGMENT_COUNT), anyLong());
+      verify(table._partitionUpsertMetadataManager, 
never()).trackSegmentForUpsertView(any());
+      verify(table._partitionUpsertMetadataManager, 
never()).trackNewlyAddedSegment(anyString());
+      if (upsertEnabled) {
+        verify(table._upsertMetadataManager).stop();
+        verify(table._upsertMetadataManager).close();
+      }
+    } finally {
+      finishConstruction.countDown();
+      additionThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  @Test(dataProvider = "upsertEnabled")
+  public void testShutdownWaitsForConsumingSegmentStart(boolean upsertEnabled)
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = Files.createTempDirectory("consuming-start").toFile();
+    RealtimeSegmentDataManager segment = 
mock(RealtimeSegmentDataManager.class);
+    MutableSegment mutableSegment = mock(MutableSegment.class);
+    when(segment.getSegment()).thenReturn(mutableSegment);
+    
when(mutableSegment.getSegmentMetadata()).thenReturn(mock(SegmentMetadata.class));
+    when(segment.decreaseReferenceCount()).thenReturn(true);
+    CountDownLatch startEntered = new CountDownLatch(1);
+    CountDownLatch finishStart = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      startEntered.countDown();
+      await(finishStart);
+      return null;
+    }).when(segment).startConsumption();
+    LifecycleTableDataManager table = new LifecycleTableDataManager(indexDir, 
segment, upsertEnabled, () -> { });
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    FutureTask<Void> shutdown = new FutureTask<>(() -> {
+      table.shutDown();
+      return null;
+    });
+    Thread additionThread = new Thread(addition, "start-consuming-segment");
+    Thread shutdownThread = new Thread(shutdown, "shutdown-consuming-table");
+    try {
+      additionThread.start();
+      await(startEntered, addition);
+      shutdownThread.start();
+      TestUtils.waitForCondition(ignored -> shutdownThread.getState() == 
Thread.State.BLOCKED || shutdown.isDone(),
+          10, 10000, "Shutdown did not reach the consuming-segment admission 
monitor");
+      assertFalse(shutdown.isDone(), "Shutdown must not offload a consumer 
before it has started");
+      verify(segment, never()).offload();
+      finishStart.countDown();
+      addition.get(10, TimeUnit.SECONDS);
+      shutdown.get(10, TimeUnit.SECONDS);
+      verify(segment).offload();
+      verify(segment).destroy();
+      assertEquals(table.getNumSegments(), 0);
+    } finally {
+      finishStart.countDown();
+      additionThread.join(10000);
+      shutdownThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  // Regression for the same-name table recreation race: a consuming add 
admitted before shutdown, paused just before
+  // it cleans up its segment data directory, must abort at the pre-cleanup 
shutdown re-check instead of resuming and
+  // deleting a directory that may already belong to a recreated same-name 
table.
+  @Test(dataProvider = "upsertEnabled")
+  public void testStaleConsumingAddAbortsBeforeDirCleanupAfterShutdown(boolean 
upsertEnabled)
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = 
Files.createTempDirectory("consuming-dir-cleanup").toFile();
+    File segmentDir = new File(indexDir, 
LifecycleTableDataManager.SEGMENT_NAME);
+    FileUtils.forceMkdir(segmentDir);
+    File marker = new File(segmentDir, "marker");
+    assertTrue(marker.createNewFile());
+    RealtimeSegmentDataManager segment = 
mock(RealtimeSegmentDataManager.class);
+    AtomicBoolean constructionRan = new AtomicBoolean();
+    CountDownLatch cleanupEntered = new CountDownLatch(1);
+    CountDownLatch finishCleanup = new CountDownLatch(1);
+    LifecycleTableDataManager table = new LifecycleTableDataManager(indexDir, 
segment, upsertEnabled,
+        () -> constructionRan.set(true), () -> {
+      cleanupEntered.countDown();
+      await(finishCleanup);
+    });
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    Thread additionThread = new Thread(addition, "cleanup-consuming-segment");
+    try {
+      additionThread.start();
+      await(cleanupEntered, addition);
+      // Shutdown must complete without waiting for the paused add.
+      table.shutDown();
+      finishCleanup.countDown();
+      ExecutionException failure = expectThrows(ExecutionException.class, () 
-> addition.get(10, TimeUnit.SECONDS));
+      assertEquals(failure.getCause().getClass(), IllegalStateException.class, 
String.valueOf(failure.getCause()));
+      // The stale add must abort before the directory cleanup and before 
constructing a segment data manager.
+      assertTrue(marker.exists(), "Stale consuming add must not delete the 
segment data directory after shutdown");
+      assertFalse(constructionRan.get(), "Stale consuming add must not 
construct a segment data manager");
+      assertEquals(table.getNumSegments(), 0);
+      verify(segment, never()).startConsumption();
+    } finally {
+      finishCleanup.countDown();
+      additionThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  // A recreated same-name table shares the data directory and the per-segment 
locks with the deleted table's manager.
+  // Its add for a colliding segment name must wait for the stale add to 
abort, find its files untouched, and complete.
+  @Test
+  public void testRecreatedTableConsumingAddSerializesBehindStaleAdd()
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = Files.createTempDirectory("consuming-recreate").toFile();
+    File segmentDir = new File(indexDir, 
LifecycleTableDataManager.SEGMENT_NAME);
+    FileUtils.forceMkdir(segmentDir);
+    File marker = new File(segmentDir, "marker");
+    assertTrue(marker.createNewFile());
+    SegmentLocks sharedLocks = new SegmentLocks();
+
+    RealtimeSegmentDataManager oldSegment = 
mock(RealtimeSegmentDataManager.class);
+    AtomicBoolean oldConstructionRan = new AtomicBoolean();
+    CountDownLatch oldCleanupEntered = new CountDownLatch(1);
+    CountDownLatch finishOldCleanup = new CountDownLatch(1);
+    LifecycleTableDataManager oldTable = new 
LifecycleTableDataManager(indexDir, oldSegment, false,
+        () -> oldConstructionRan.set(true), () -> {
+      oldCleanupEntered.countDown();
+      await(finishOldCleanup);
+    }, sharedLocks);
+
+    RealtimeSegmentDataManager newSegment = 
mock(RealtimeSegmentDataManager.class);
+    MutableSegment mutableSegment = mock(MutableSegment.class);
+    when(newSegment.getSegment()).thenReturn(mutableSegment);
+    
when(mutableSegment.getSegmentMetadata()).thenReturn(mock(SegmentMetadata.class));
+    when(newSegment.decreaseReferenceCount()).thenReturn(true);
+    AtomicBoolean newConstructionRan = new AtomicBoolean();
+    CountDownLatch newCleanupEntered = new CountDownLatch(1);
+    CountDownLatch finishNewCleanup = new CountDownLatch(1);
+
+    FutureTask<Void> oldAddition = new FutureTask<>(() -> {
+      oldTable.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    Thread oldAdditionThread = new Thread(oldAddition, "stale-consuming-add");
+    Thread newAdditionThread = null;
+    try {
+      oldAdditionThread.start();
+      await(oldCleanupEntered, oldAddition);
+      // Table deletion: shutdown completes while the stale add is paused 
inside the segment lock.
+      oldTable.shutDown();
+
+      LifecycleTableDataManager newTable = new 
LifecycleTableDataManager(indexDir, newSegment, false,
+          () -> newConstructionRan.set(true), () -> {
+        newCleanupEntered.countDown();
+        await(finishNewCleanup);
+      }, sharedLocks);
+      FutureTask<Void> newAddition = new FutureTask<>(() -> {
+        newTable.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+        return null;
+      });
+      newAdditionThread = new Thread(newAddition, "recreated-consuming-add");
+      newAdditionThread.start();
+      // The recreated table's add must block on the shared per-segment lock 
while the stale add is in flight.
+      Thread newThread = newAdditionThread;
+      TestUtils.waitForCondition(ignored -> newThread.getState() == 
Thread.State.WAITING || newAddition.isDone(),
+          10, 10000, "Recreated table's add did not block on the shared 
segment lock");
+      assertFalse(newAddition.isDone(), "Recreated table's add must wait for 
the stale add to finish");
+      assertEquals(newCleanupEntered.getCount(), 1,
+          "Recreated table's add must not enter the critical section while the 
stale add holds the segment lock");
+
+      finishOldCleanup.countDown();
+      ExecutionException failure = expectThrows(ExecutionException.class, () 
-> oldAddition.get(10, TimeUnit.SECONDS));
+      assertEquals(failure.getCause().getClass(), IllegalStateException.class, 
String.valueOf(failure.getCause()));
+      assertFalse(oldConstructionRan.get(), "Stale consuming add must not 
construct a segment data manager");
+
+      // The recreated table's add now proceeds; the stale add must have left 
the segment directory untouched.
+      await(newCleanupEntered, newAddition);
+      assertTrue(marker.exists(), "Stale consuming add must not delete the 
recreated table's segment directory");
+      finishNewCleanup.countDown();
+      newAddition.get(10, TimeUnit.SECONDS);
+      assertTrue(newConstructionRan.get());
+      assertEquals(newTable.getNumSegments(), 1);
+      verify(newSegment).startConsumption();
+      assertFalse(marker.exists(), "Recreated table's add owns the directory 
cleanup");
+    } finally {
+      finishOldCleanup.countDown();
+      finishNewCleanup.countDown();
+      oldAdditionThread.join(10000);
+      if (newAdditionThread != null) {
+        newAdditionThread.join(10000);
+      }
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  @Test
+  public void 
testShutdownDoesNotDestroyConsumingSegmentDuringOnlineTransition()
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = 
Files.createTempDirectory("online-transition-shutdown").toFile();
+    RealtimeSegmentDataManager segment = 
mock(RealtimeSegmentDataManager.class);
+    // Real reference counting on the mock: shutdown's release must not be the 
last one while the transition runs.
+    FieldUtils.writeField(segment, "_referenceCount", 1, true);
+    doCallRealMethod().when(segment).increaseReferenceCount();
+    doCallRealMethod().when(segment).decreaseReferenceCount();
+    doCallRealMethod().when(segment).getReferenceCount();
+    MutableSegment mutableSegment = mock(MutableSegment.class);
+    
when(mutableSegment.getSegmentMetadata()).thenReturn(mock(SegmentMetadata.class));
+    when(segment.getSegment()).thenReturn(mutableSegment);
+    LifecycleTableDataManager table = spy(new 
LifecycleTableDataManager(indexDir, segment, false, () -> { }));
+    CountDownLatch transitionStarted = new CountDownLatch(1);
+    CountDownLatch finishTransition = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      transitionStarted.countDown();
+      await(finishTransition);
+      return null;
+    }).when(segment).goOnlineFromConsuming(any());
+    table.addConsumingSegment(LifecycleTableDataManager.SEGMENT_NAME);
+    assertEquals(table.getNumSegments(), 1);
+    assertEquals(segment.getReferenceCount(), 1);
+    SegmentZKMetadata committedMetadata = new 
SegmentZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    committedMetadata.setStatus(CommonConstants.Segment.Realtime.Status.DONE);
+    
doReturn(committedMetadata).when(table).fetchZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    FutureTask<Void> transition = new FutureTask<>(() -> {
+      table.addOnlineSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    Thread transitionThread = new Thread(transition, "consuming-to-online");
+    FutureTask<Void> shutdown = new FutureTask<>(() -> {
+      table.shutDown();
+      return null;
+    });
+    Thread shutdownThread = new Thread(shutdown, "table-shutdown");
+    try {
+      transitionThread.start();
+      await(transitionStarted, transition);
+      assertEquals(segment.getReferenceCount(), 2, "Transition holds its own 
reference");
+      shutdownThread.start();
+      // Shutdown offloads and releases the segment without waiting for the 
transition, and must not destroy it.
+      shutdown.get(10, TimeUnit.SECONDS);
+      verify(segment).offload();
+      verify(segment, never()).destroy();
+      assertEquals(table.getNumSegments(), 0);
+      assertEquals(segment.getReferenceCount(), 1, "Shutdown released the 
map's reference only");
+      finishTransition.countDown();
+      transition.get(10, TimeUnit.SECONDS);
+      verify(segment).destroy();
+      assertEquals(segment.getReferenceCount(), 0);
+    } finally {
+      finishTransition.countDown();
+      transitionThread.join(10000);
+      shutdownThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  // Without the shutdown re-check the wait would run to the 10-minute default 
download timeout; fail fast instead.
+  @Test(timeOut = 60_000)
+  public void testCommittingSegmentDownloadAbortsAfterShutdown()
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = 
Files.createTempDirectory("committing-download-shutdown").toFile();
+    LifecycleTableDataManager table = spy(
+        new LifecycleTableDataManager(indexDir, 
mock(RealtimeSegmentDataManager.class), false, () -> { }));
+    // The pauseless wait derives its timeout from the stream config, so give 
the cached table config one.
+    Map<String, String> streamConfigs = new HashMap<>();
+    streamConfigs.put("streamType", "kafka");
+    streamConfigs.put("stream.kafka.topic.name", "lifecycle");
+    TableConfig tableConfig = new 
TableConfigBuilder(TableType.REALTIME).setTableName("lifecycle")
+        .setStreamConfigs(streamConfigs).build();
+    FieldUtils.writeField(table, "_cachedTableConfigAndSchema",
+        Pair.of(tableConfig, 
table.getCachedTableConfigAndSchema().getRight()), true);
+    SegmentZKMetadata committing = new 
SegmentZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    committing.setStatus(CommonConstants.Segment.Realtime.Status.COMMITTING);
+    
doReturn(committing).when(table).fetchZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    try {
+      table.shutDown();
+      // A stale wait must abort before polling ZK again instead of holding 
the shared segment lock until the timeout.
+      IllegalStateException failure =
+          expectThrows(IllegalStateException.class, () -> 
table.downloadSegment(committing));
+      assertTrue(failure.getMessage().contains("already shut down"), 
failure.getMessage());
+      verify(table, 
never()).fetchZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    } finally {
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  @DataProvider
+  public Object[][] upsertConsistencyModes() {
+    return new Object[][]{{ConsistencyMode.NONE}, {ConsistencyMode.SYNC}, 
{ConsistencyMode.SNAPSHOT}};
+  }
+
+  @Test(dataProvider = "upsertEnabled")
+  public void testOnlinePreloadDoesNotCreatePartitionAfterShutdown(boolean 
upsertEnabled)
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = 
Files.createTempDirectory("online-preload-shutdown").toFile();
+    LifecycleTableDataManager table = spy(
+        new LifecycleTableDataManager(indexDir, 
mock(RealtimeSegmentDataManager.class), upsertEnabled, () -> { }));
+    TableDedupMetadataManager dedup = mock(TableDedupMetadataManager.class);
+    PartitionDedupMetadataManager partitionDedup = 
mock(PartitionDedupMetadataManager.class);
+    if (upsertEnabled) {
+      
when(table._upsertMetadataManager.getContext().isPreloadEnabled()).thenReturn(true);
+    } else {
+      FieldUtils.writeField(table, "_tableDedupMetadataManager", dedup, true);
+      DedupContext context = mock(DedupContext.class);
+      when(dedup.getContext()).thenReturn(context);
+      when(context.isPreloadEnabled()).thenReturn(true);
+      when(dedup.getOrCreatePartitionManager(0)).thenReturn(partitionDedup);
+    }
+    SegmentZKMetadata metadata = new 
SegmentZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    metadata.setStatus(CommonConstants.Segment.Realtime.Status.DONE);
+    
doReturn(metadata).when(table).fetchZKMetadata(LifecycleTableDataManager.SEGMENT_NAME);
+    CountDownLatch loadingConfig = new CountDownLatch(1);
+    CountDownLatch finishLoadingConfig = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      loadingConfig.countDown();
+      await(finishLoadingConfig);
+      return invocation.callRealMethod();
+    }).when(table).fetchIndexLoadingConfig();
+    FutureTask<Void> addition = new FutureTask<>(() -> {
+      table.addOnlineSegment(LifecycleTableDataManager.SEGMENT_NAME);
+      return null;
+    });
+    Thread additionThread = new Thread(addition, "preload-online-segment");
+    try {
+      additionThread.start();
+      await(loadingConfig, addition);
+      table.shutDown();
+      finishLoadingConfig.countDown();
+      ExecutionException failure = expectThrows(ExecutionException.class, () 
-> addition.get(10, TimeUnit.SECONDS));
+      verify(table._upsertMetadataManager, 
never()).getOrCreatePartitionManager(0);
+      verify(table._partitionUpsertMetadataManager, 
never()).preloadSegments(any());
+      verify(dedup, never()).getOrCreatePartitionManager(0);
+      verify(partitionDedup, never()).preloadSegments(any());
+      assertTrue(failure.getCause() instanceof IllegalStateException);
+      assertEquals(table.getNumSegments(), 0);
+    } finally {
+      finishLoadingConfig.countDown();
+      additionThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  @Test(dataProvider = "upsertConsistencyModes")
+  public void testShutdownWaitsForUpsertSegmentReplacement(ConsistencyMode 
consistencyMode)
+      throws Exception {
+    ServerMetrics.register(mock(ServerMetrics.class));
+    File indexDir = 
Files.createTempDirectory("upsert-replacement-shutdown").toFile();
+    LifecycleTableDataManager table =
+        new LifecycleTableDataManager(indexDir, 
mock(RealtimeSegmentDataManager.class), true, () -> { });
+    
when(table._upsertMetadataManager.getContext().getConsistencyMode()).thenReturn(consistencyMode);
+    ImmutableSegment oldSegment = 
immutableSegment(LifecycleTableDataManager.SEGMENT_NAME);
+    ImmutableSegment newSegment = 
immutableSegment(LifecycleTableDataManager.SEGMENT_NAME);
+    table.addSegment(oldSegment, null);
+    CountDownLatch replacementEntered = new CountDownLatch(1);
+    CountDownLatch finishReplacement = new CountDownLatch(1);
+    doAnswer(invocation -> {
+      replacementEntered.countDown();
+      await(finishReplacement);
+      return null;
+    }).when(table._partitionUpsertMetadataManager).replaceSegment(newSegment, 
oldSegment);
+    FutureTask<Void> replacement = new FutureTask<>(() -> {
+      table.addSegment(newSegment, null);
+      return null;
+    });
+    FutureTask<Void> shutdown = new FutureTask<>(() -> {
+      table.shutDown();
+      return null;
+    });
+    Thread replacementThread = new Thread(replacement, 
"replace-upsert-segment");
+    Thread shutdownThread = new Thread(shutdown, "shutdown-upsert-table");
+    try {
+      replacementThread.start();
+      await(replacementEntered, replacement);
+      
assertEquals(table.getSegmentDataManager(LifecycleTableDataManager.SEGMENT_NAME).hasMultiSegments(),
+          consistencyMode != ConsistencyMode.NONE);
+      shutdownThread.start();
+      TestUtils.waitForCondition(ignored -> table.isShutDown()
+              && (shutdownThread.getState() == Thread.State.WAITING || 
shutdown.isDone()), 10, 10000,
+          "Shutdown did not wait for the admitted upsert replacement");
+      assertFalse(shutdown.isDone(), "Shutdown must wait for the complete 
upsert replacement");
+      verify(table._upsertMetadataManager, never()).stop();
+      verify(oldSegment, never()).destroy();
+      verify(newSegment, never()).destroy();
+      finishReplacement.countDown();
+      replacement.get(10, TimeUnit.SECONDS);
+      shutdown.get(10, TimeUnit.SECONDS);
+      verify(oldSegment).offload();
+      verify(oldSegment).destroy();
+      verify(newSegment).offload();
+      verify(newSegment).destroy();
+      verify(table._upsertMetadataManager).stop();
+      verify(table._upsertMetadataManager).close();
+      assertEquals(table.getNumSegments(), 0);
+
+      ImmutableSegment lateSegment = immutableSegment("lifecycle__1__0__1");
+      expectThrows(IllegalStateException.class, () -> 
table.addSegment(lateSegment, null));
+      verify(lateSegment).destroy();
+      verify(table._upsertMetadataManager, 
never()).getOrCreatePartitionManager(1);
+      assertEquals(table.getNumSegments(), 0);
+    } finally {
+      finishReplacement.countDown();
+      replacementThread.join(10000);
+      shutdownThread.join(10000);
+      FileUtils.deleteDirectory(indexDir);
+    }
+  }
+
+  private static ImmutableSegment immutableSegment(String segmentName) {
+    ImmutableSegment segment = mock(ImmutableSegment.class);
+    SegmentMetadata metadata = mock(SegmentMetadata.class);
+    when(segment.getSegmentName()).thenReturn(segmentName);
+    when(segment.getSegmentMetadata()).thenReturn(metadata);
+    when(metadata.getName()).thenReturn(segmentName);
+    return segment;
+  }
+
+  private static void await(CountDownLatch latch, FutureTask<?> operation)
+      throws Exception {
+    TestUtils.waitForCondition(ignored -> latch.getCount() == 0 || 
operation.isDone(), 10, 10000,
+        "Consuming segment operation did not reach the lifecycle barrier");
+    if (operation.isDone()) {
+      operation.get();
+    }
+    assertEquals(latch.getCount(), 0L);
+  }
+
+  private static void await(CountDownLatch latch) {
+    try {
+      assertTrue(latch.await(10, TimeUnit.SECONDS), "Timed out waiting for 
consuming segment lifecycle operation");
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+      throw new AssertionError(e);
+    }
+  }
+
+  /// Exercises public segment admission and shutdown while controlling 
construction, startup and upsert replacement.
+  private static class LifecycleTableDataManager extends 
RealtimeTableDataManager {
+    private static final String SEGMENT_NAME = "lifecycle__0__0__1";
+    private final TableConfig _tableConfig =
+        new 
TableConfigBuilder(TableType.REALTIME).setTableName("lifecycle").build();
+    private final Schema _schema = new 
Schema.SchemaBuilder().setSchemaName("lifecycle")
+        .addSingleValueDimension("value", DataType.STRING).build();
+    private final RealtimeSegmentDataManager _segment;
+    private final Runnable _beforeConstructionReturns;
+    private final Runnable _beforeAdmissionRecheck;
+    private final PartitionUpsertMetadataManager 
_partitionUpsertMetadataManager =
+        mock(PartitionUpsertMetadataManager.class);
+    private final TableUpsertMetadataManager _upsertMetadataManager = 
mock(TableUpsertMetadataManager.class);
+
+    LifecycleTableDataManager(File indexDir, RealtimeSegmentDataManager 
segment, boolean upsertEnabled,
+        Runnable beforeConstructionReturns)
+        throws IllegalAccessException {
+      this(indexDir, segment, upsertEnabled, beforeConstructionReturns, () -> 
{ }, new SegmentLocks());
+    }
+
+    LifecycleTableDataManager(File indexDir, RealtimeSegmentDataManager 
segment, boolean upsertEnabled,
+        Runnable beforeConstructionReturns, Runnable beforeAdmissionRecheck)
+        throws IllegalAccessException {
+      this(indexDir, segment, upsertEnabled, beforeConstructionReturns, 
beforeAdmissionRecheck, new SegmentLocks());
+    }
+
+    LifecycleTableDataManager(File indexDir, RealtimeSegmentDataManager 
segment, boolean upsertEnabled,
+        Runnable beforeConstructionReturns, Runnable beforeAdmissionRecheck, 
SegmentLocks segmentLocks)
+        throws IllegalAccessException {
+      super(null);
+      _indexDir = indexDir;
+      _tableNameWithType = _tableConfig.getTableName();
+      _cachedTableConfigAndSchema = Pair.of(_tableConfig, _schema);
+      _logger = LoggerFactory.getLogger(LifecycleTableDataManager.class);
+      _segmentLocks = segmentLocks;
+      _recentlyDeletedSegments = CacheBuilder.newBuilder().build();
+      _segment = segment;
+      _beforeConstructionReturns = beforeConstructionReturns;
+      _beforeAdmissionRecheck = beforeAdmissionRecheck;
+      FieldUtils.writeField(this, "_ingestionDelayTracker", 
mock(IngestionDelayTracker.class), true);
+      if (upsertEnabled) {
+        _tableUpsertMetadataManager = _upsertMetadataManager;
+        
when(_tableUpsertMetadataManager.getContext()).thenReturn(mock(UpsertContext.class));
+        
when(_tableUpsertMetadataManager.getOrCreatePartitionManager(0)).thenReturn(_partitionUpsertMetadataManager);
+      }
+    }
+
+    @Override
+    void beforeConsumingSegmentAdmissionRecheck(String segmentName) {
+      _beforeAdmissionRecheck.run();
+    }
+
+    @Override
+    public SegmentZKMetadata fetchZKMetadata(String segmentName) {
+      SegmentZKMetadata metadata = new SegmentZKMetadata(segmentName);
+      metadata.setStatus(CommonConstants.Segment.Realtime.Status.IN_PROGRESS);
+      return metadata;
+    }
+
+    @Override
+    public IndexLoadingConfig fetchIndexLoadingConfig() {
+      return new IndexLoadingConfig(_tableConfig, _schema);
+    }
+
+    /// Ordinary ONLINE and CONSUMING loads resolve the cached config. Route 
it through [#fetchIndexLoadingConfig()]
+    /// so tests can pause a load at the point where its config is resolved, 
before any preload or admission.
+    @Override
+    protected IndexLoadingConfig getCachedIndexLoadingConfig() {
+      return fetchIndexLoadingConfig();
+    }
+
+    @Override
+    protected RealtimeSegmentDataManager 
createRealtimeSegmentDataManager(SegmentZKMetadata zkMetadata,
+        TableConfig tableConfig, IndexLoadingConfig indexLoadingConfig, Schema 
schema, LLCSegmentName llcSegmentName,
+        ConsumerCoordinator consumerCoordinator, 
PartitionUpsertMetadataManager partitionUpsertMetadataManager,
+        PartitionDedupMetadataManager partitionDedupMetadataManager, 
BooleanSupplier isTableReadyToConsumeData) {
+      _beforeConstructionReturns.run();
+      return _segment;
+    }
+  }
 
   @Test
   public void testSetDefaultTimeValueIfInvalid() {
diff --git 
a/pinot-server/src/main/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManager.java
 
b/pinot-server/src/main/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManager.java
index 02fadfae53f..99c23eaa6c6 100644
--- 
a/pinot-server/src/main/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManager.java
+++ 
b/pinot-server/src/main/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManager.java
@@ -22,6 +22,8 @@ import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Preconditions;
 import com.google.common.cache.Cache;
 import com.google.common.cache.CacheBuilder;
+import com.google.common.cache.CacheLoader;
+import com.google.common.cache.LoadingCache;
 import com.google.common.util.concurrent.ThreadFactoryBuilder;
 import java.io.File;
 import java.io.IOException;
@@ -33,8 +35,8 @@ import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
 import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
 import java.util.function.BooleanSupplier;
 import javax.annotation.Nullable;
 import javax.annotation.concurrent.ThreadSafe;
@@ -87,11 +89,20 @@ public class HelixInstanceDataManager implements 
InstanceDataManager {
   private static final Logger LOGGER = 
LoggerFactory.getLogger(HelixInstanceDataManager.class);
 
   private final Map<String, TableDataManager> _tableDataManagerMap = new 
ConcurrentHashMap<>();
+  /// Serializes table creation with the previous owner's entire shutdown, 
including table-scoped resource cleanup
+  /// such as the segment build time lease extender. Acquire it before 
changing [#_tableDataManagerMap], and release
+  /// it before calling the table's segment-add methods. Values are weak so 
idle tables do not accumulate locks; a
+  /// lock stays strongly reachable from its holder for as long as it is held.
+  private final LoadingCache<String, Lock> _tableLifecycleLocks =
+      CacheBuilder.newBuilder().weakValues().build(CacheLoader.from(() -> new 
ReentrantLock()));
 
   // Logical table metadata cache to cache logical table configs, schemas, and 
offline/realtime table configs.
   private final LogicalTableMetadataCache _logicalTableMetadataCache = new 
LogicalTableMetadataCache();
 
-  // TODO: Consider making segment locks per table instead of per instance
+  /// Intentionally shared across all table data managers, including 
successive incarnations of the same table name:
+  /// a deleted table's stale segment operation and a recreated same-name 
table's operation on a colliding segment
+  /// name must serialize on one lock. The shutdown re-checks in 
`BaseTableDataManager.moveSegment` and
+  /// `RealtimeTableDataManager.doAddConsumingSegment` rely on this, so do not 
make these locks per table.
   private final SegmentLocks _segmentLocks = new SegmentLocks();
 
   private HelixInstanceDataManagerConfig _instanceDataManagerConfig;
@@ -308,37 +319,62 @@ public class HelixInstanceDataManager implements 
InstanceDataManager {
   @Override
   public void deleteTable(String tableNameWithType, long deletionTimeMs)
       throws Exception {
-    AtomicReference<TableDataManager> tableDataManagerRef = new 
AtomicReference<>();
-    _tableDataManagerMap.computeIfPresent(tableNameWithType, (k, v) -> {
-      _recentlyDeletedTables.put(k, deletionTimeMs);
-      tableDataManagerRef.set(v);
-      return null;
-    });
-    TableDataManager tableDataManager = tableDataManagerRef.get();
-    if (tableDataManager == null) {
-      LOGGER.warn("Failed to find table data manager for table: {}, skip 
deleting the table", tableNameWithType);
-      return;
+    Lock lifecycleLock = _tableLifecycleLocks.getUnchecked(tableNameWithType);
+    lifecycleLock.lock();
+    try {
+      // The first segment callback can arrive after deletion, before a 
manager has ever been created.
+      // Keep the newest deletion timestamp when duplicate or out-of-order 
deletion messages arrive.
+      _recentlyDeletedTables.asMap().merge(tableNameWithType, deletionTimeMs, 
Math::max);
+      TableDataManager tableDataManager = 
_tableDataManagerMap.remove(tableNameWithType);
+      if (tableDataManager == null) {
+        LOGGER.warn("Failed to find table data manager for table: {}, skip 
deleting the table", tableNameWithType);
+        return;
+      }
+      // Shutdown can wait for segment callbacks, so do not run it inside a 
map computation.
+      LOGGER.info("Shutting down table data manager for table: {}", 
tableNameWithType);
+      tableDataManager.setDeleted(true);
+      tableDataManager.shutDown();
+      LOGGER.info("Finished shutting down table data manager for table: {}", 
tableNameWithType);
+    } finally {
+      lifecycleLock.unlock();
     }
-    LOGGER.info("Shutting down table data manager for table: {}", 
tableNameWithType);
-    tableDataManager.setDeleted(true);
-    tableDataManager.shutDown();
-    LOGGER.info("Finished shutting down table data manager for table: {}", 
tableNameWithType);
   }
 
   @Override
   public void addOnlineSegment(String tableNameWithType, String segmentName)
       throws Exception {
-    _tableDataManagerMap.computeIfAbsent(tableNameWithType, 
this::createTableDataManager).addOnlineSegment(segmentName);
+    
getOrCreateTableDataManager(tableNameWithType).addOnlineSegment(segmentName);
   }
 
   @Override
   public void addConsumingSegment(String realtimeTableName, String segmentName)
       throws Exception {
-    _tableDataManagerMap.computeIfAbsent(realtimeTableName, 
this::createTableDataManager)
-        .addConsumingSegment(segmentName);
+    
getOrCreateTableDataManager(realtimeTableName).addConsumingSegment(segmentName);
   }
 
-  private TableDataManager createTableDataManager(String tableNameWithType) {
+  /// Returns the table data manager, creating and starting it under the 
table's lifecycle lock when absent. The
+  /// lock-free fast path can return a manager that a concurrent 
[#deleteTable] is shutting down; that manager's
+  /// segment-add methods reject the operation once shutdown has closed 
admission. The slow path waits for such a
+  /// shutdown to finish before creating the replacement.
+  private TableDataManager getOrCreateTableDataManager(String 
tableNameWithType) {
+    TableDataManager tableDataManager = 
_tableDataManagerMap.get(tableNameWithType);
+    if (tableDataManager != null) {
+      return tableDataManager;
+    }
+    Lock lifecycleLock = _tableLifecycleLocks.getUnchecked(tableNameWithType);
+    lifecycleLock.lock();
+    try {
+      return _tableDataManagerMap.computeIfAbsent(tableNameWithType, 
this::createTableDataManager);
+    } finally {
+      lifecycleLock.unlock();
+    }
+  }
+
+  /// Creates and starts a table data manager; callers must hold the table's 
lifecycle lock. A recently deleted table
+  /// is recreated only from a table config created after its newest recorded 
deletion, and the deletion record is
+  /// cleared only after the manager has started, so a failed creation keeps 
rejecting stale configs.
+  @VisibleForTesting
+  TableDataManager createTableDataManager(String tableNameWithType) {
     LOGGER.info("Creating table data manager for table: {}", 
tableNameWithType);
     TableConfig tableConfig;
     Long tableDeleteTimeMs = 
_recentlyDeletedTables.getIfPresent(tableNameWithType);
@@ -353,9 +389,8 @@ public class HelixInstanceDataManager implements 
InstanceDataManager {
       tableConfig = tableConfigAndStat.getLeft();
       long tableCreationTimeMs = tableConfigAndStat.getRight().getCtime();
       Preconditions.checkState(tableCreationTimeMs > tableDeleteTimeMs,
-          "Table: %s was recently deleted (deleted %dms ago) but the table 
config was created before that (created "
-              + "%dms ago)", tableNameWithType, currentTimeMs - 
tableDeleteTimeMs, currentTimeMs - tableCreationTimeMs);
-      _recentlyDeletedTables.invalidate(tableNameWithType);
+          "Table: %s was recently deleted (deleted %sms ago) but the table 
config was created before that (created "
+              + "%sms ago)", tableNameWithType, currentTimeMs - 
tableDeleteTimeMs, currentTimeMs - tableCreationTimeMs);
     } else {
       tableConfig = ZKMetadataProvider.getTableConfig(_propertyStore, 
tableNameWithType);
       Preconditions.checkState(tableConfig != null, "Failed to find table 
config for table: %s", tableNameWithType);
@@ -369,6 +404,7 @@ public class HelixInstanceDataManager implements 
InstanceDataManager {
             _isServerReadyToServeQueries, 
_serverIngestionOomProtectionThrottleState, _enableAsyncSegmentRefresh,
             _reloadJobStatusCache);
     tableDataManager.start();
+    _recentlyDeletedTables.invalidate(tableNameWithType);
     LOGGER.info("Created table data manager for table: {}", tableNameWithType);
     return tableDataManager;
   }
diff --git 
a/pinot-server/src/test/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManagerLifecycleTest.java
 
b/pinot-server/src/test/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManagerLifecycleTest.java
new file mode 100644
index 00000000000..aa03e394be8
--- /dev/null
+++ 
b/pinot-server/src/test/java/org/apache/pinot/server/starter/helix/HelixInstanceDataManagerLifecycleTest.java
@@ -0,0 +1,380 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.server.starter.helix;
+
+import com.google.common.cache.CacheBuilder;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Function;
+import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.helix.AccessOption;
+import org.apache.helix.store.zk.ZkHelixPropertyStore;
+import org.apache.helix.zookeeper.datamodel.ZNRecord;
+import org.apache.pinot.common.metadata.ZKMetadataProvider;
+import org.apache.pinot.common.utils.config.SchemaSerDeUtils;
+import org.apache.pinot.common.utils.config.TableConfigSerDeUtils;
+import org.apache.pinot.core.data.manager.provider.TableDataManagerProvider;
+import 
org.apache.pinot.core.data.manager.realtime.SegmentBuildTimeLeaseExtender;
+import org.apache.pinot.segment.local.data.manager.TableDataManager;
+import org.apache.pinot.spi.config.table.TableType;
+import org.apache.pinot.spi.data.Schema;
+import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
+import org.apache.zookeeper.data.Stat;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotSame;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertSame;
+import static org.testng.Assert.assertThrows;
+import static org.testng.Assert.assertTrue;
+
+
+/// Exercises table lifecycle ordering with controlled shutdown and segment 
callbacks, without starting a server.
+public class HelixInstanceDataManagerLifecycleTest {
+  private static final long TIMEOUT_SECONDS = 10;
+
+  @DataProvider
+  public Object[][] segmentTypes() {
+    return new Object[][]{{false}, {true}};
+  }
+
+  @DataProvider
+  public Object[][] tableManagerPresence() {
+    return new Object[][]{{false}, {true}};
+  }
+
+  @Test(dataProvider = "tableManagerPresence")
+  @SuppressWarnings("unchecked")
+  public void testDeletedTableCreationChecksConfigTimestamp(boolean 
existingManager)
+      throws Exception {
+    String table = "deleted_REALTIME";
+    ZNRecord configRecord = TableConfigSerDeUtils.toZNRecord(
+        new 
TableConfigBuilder(TableType.REALTIME).setTableName("deleted").build());
+    AtomicReference<ZNRecord> currentConfig = new 
AtomicReference<>(configRecord);
+    AtomicLong creationTimeMs = new AtomicLong(100L);
+    ZkHelixPropertyStore<ZNRecord> propertyStore = 
mock(ZkHelixPropertyStore.class);
+    
when(propertyStore.get(eq(ZKMetadataProvider.constructPropertyStorePathForResourceConfig(table)),
 any(),
+        eq(AccessOption.PERSISTENT))).thenAnswer(invocation -> {
+          Stat stat = invocation.getArgument(1);
+          if (stat != null) {
+            stat.setCtime(creationTimeMs.get());
+          }
+          return currentConfig.get();
+        });
+    
when(propertyStore.get(eq(ZKMetadataProvider.constructPropertyStorePathForSchema("deleted")),
 any(),
+        eq(AccessOption.PERSISTENT))).thenReturn(
+            SchemaSerDeUtils.toZNRecord(new 
Schema.SchemaBuilder().setSchemaName("deleted").build()));
+    TableDataManager oldManager = mock(TableDataManager.class);
+    TableDataManager replacement = mock(TableDataManager.class);
+    TableDataManagerProvider provider = mock(TableDataManagerProvider.class);
+    when(provider.getTableDataManager(any(), any(), any(), any(), any(), 
any(), any(), any(), any(), anyBoolean(),
+        any()))
+        .thenReturn(existingManager ? oldManager : replacement, replacement);
+    HelixInstanceDataManager instance = new HelixInstanceDataManager();
+    instance._recentlyDeletedTables = CacheBuilder.newBuilder().build();
+    FieldUtils.writeField(instance, "_propertyStore", propertyStore, true);
+    FieldUtils.writeField(instance, "_tableDataManagerProvider", provider, 
true);
+    if (existingManager) {
+      instance.addConsumingSegment(table, "old");
+    }
+
+    instance.deleteTable(table, 200L);
+    instance.deleteTable(table, 200L);
+    instance.deleteTable(table, 50L);
+    assertThrows(IllegalStateException.class, () -> 
instance.addConsumingSegment(table, "stale"));
+    currentConfig.set(null);
+    assertThrows(IllegalStateException.class, () -> 
instance.addConsumingSegment(table, "missing"));
+    assertNull(instance.getTableDataManager(table));
+    verify(provider, times(existingManager ? 1 : 
0)).getTableDataManager(any(), any(), any(), any(), any(), any(),
+        any(), any(), any(), anyBoolean(), any());
+
+    // A genuinely recreated table is allowed, but a failed start must retain 
the deletion guard.
+    currentConfig.set(configRecord);
+    creationTimeMs.set(201L);
+    doThrow(new IllegalStateException("start 
failure")).doNothing().when(replacement).start();
+    assertThrows(IllegalStateException.class, () -> 
instance.addConsumingSegment(table, "failed"));
+    creationTimeMs.set(100L);
+    assertThrows(IllegalStateException.class, () -> 
instance.addConsumingSegment(table, "stale-after-failure"));
+    assertNull(instance.getTableDataManager(table));
+    creationTimeMs.set(201L);
+    instance.addConsumingSegment(table, "new");
+    assertSame(instance.getTableDataManager(table), replacement);
+    assertNull(instance._recentlyDeletedTables.getIfPresent(table));
+    verify(replacement).addConsumingSegment("new");
+  }
+
+  @Test(dataProvider = "segmentTypes")
+  public void testRecreationWaitsForOldLeaseShutdown(boolean consuming)
+      throws Exception {
+    String table = "lifecycle_" + consuming + "_REALTIME";
+    TableDataManager oldManager = mock(TableDataManager.class);
+    TableDataManager replacement = mock(TableDataManager.class);
+    CountDownLatch shutdownEntered = new CountDownLatch(1);
+    CountDownLatch finishShutdown = new CountDownLatch(1);
+    CountDownLatch shutdownFinished = new CountDownLatch(1);
+    AtomicInteger creations = new AtomicInteger();
+    AtomicReference<SegmentBuildTimeLeaseExtender> replacementLease = new 
AtomicReference<>();
+    SegmentBuildTimeLeaseExtender oldLease = 
SegmentBuildTimeLeaseExtender.getOrCreate("server", null, table);
+    TestInstanceDataManager instance = new TestInstanceDataManager(name -> {
+      if (creations.getAndIncrement() == 0) {
+        return oldManager;
+      }
+      assertTrue(shutdownFinished.getCount() == 0, "Replacement initialized 
before old lease shutdown completed");
+      replacementLease.set(SegmentBuildTimeLeaseExtender.getOrCreate("server", 
null, name));
+      return replacement;
+    });
+    doAnswer(invocation -> {
+      shutdownEntered.countDown();
+      await(finishShutdown);
+      oldLease.shutDown();
+      shutdownFinished.countDown();
+      return null;
+    }).when(oldManager).shutDown();
+    addSegment(instance, table, "old", consuming);
+    FutureTask<Void> deletion = new FutureTask<>(() -> {
+      instance.deleteTable(table, 1L);
+      return null;
+    });
+    Thread deletionThread = new Thread(deletion, "delete-table");
+    FutureTask<Void> recreation = new FutureTask<>(() -> {
+      addSegment(instance, table, "new", consuming);
+      return null;
+    });
+    Thread recreationThread = new Thread(recreation, "recreate-table");
+    try {
+      deletionThread.start();
+      await(shutdownEntered);
+      recreationThread.start();
+      // Wait for the competing operation to reach the lock, or expose an 
incorrectly completed creation.
+      long deadline = System.nanoTime() + 
TimeUnit.SECONDS.toNanos(TIMEOUT_SECONDS);
+      while (recreationThread.getState() != Thread.State.WAITING && 
!recreation.isDone()
+          && System.nanoTime() < deadline) {
+        Thread.sleep(10L);
+      }
+      assertFalse(recreation.isDone(), "Recreation must wait for the old 
table's shutdown");
+      assertSame(recreationThread.getState(), Thread.State.WAITING, 
"Recreation did not reach the lifecycle lock");
+      finishShutdown.countDown();
+      deletion.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      recreation.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      assertNotSame(replacementLease.get(), oldLease);
+      assertSame(SegmentBuildTimeLeaseExtender.getLeaseExtender(table), 
replacementLease.get());
+      assertSame(instance.getTableDataManager(table), replacement);
+    } finally {
+      finishShutdown.countDown();
+      join(deletionThread);
+      join(recreationThread);
+      SegmentBuildTimeLeaseExtender lease = 
SegmentBuildTimeLeaseExtender.getLeaseExtender(table);
+      if (lease != null) {
+        lease.shutDown();
+      }
+    }
+  }
+
+  @Test
+  public void testOtherTableCanInitializeDuringShutdown()
+      throws Exception {
+    String table = "Aa_REALTIME";
+    String otherTable = "BB_REALTIME";
+    assertEquals(table.hashCode(), otherTable.hashCode());
+    TableDataManager oldManager = mock(TableDataManager.class);
+    TableDataManager otherManager = mock(TableDataManager.class);
+    CountDownLatch shutdownEntered = new CountDownLatch(1);
+    CountDownLatch finishShutdown = new CountDownLatch(1);
+    TestInstanceDataManager instance = new TestInstanceDataManager(name -> 
name.equals(table)
+        ? oldManager
+        : otherManager);
+    doAnswer(invocation -> {
+      shutdownEntered.countDown();
+      await(finishShutdown);
+      return null;
+    }).when(oldManager).shutDown();
+    instance.addConsumingSegment(table, "old");
+    FutureTask<Void> deletion = new FutureTask<>(() -> {
+      instance.deleteTable(table, 1L);
+      return null;
+    });
+    Thread deletionThread = new Thread(deletion, "delete-other-table");
+    FutureTask<Void> creation = new FutureTask<>(() -> {
+      instance.addConsumingSegment(otherTable, "other");
+      return null;
+    });
+    Thread creationThread = new Thread(creation, "create-other-table");
+    try {
+      deletionThread.start();
+      await(shutdownEntered);
+      creationThread.start();
+      creation.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      verify(otherManager).addConsumingSegment("other");
+    } finally {
+      finishShutdown.countDown();
+      join(deletionThread);
+      join(creationThread);
+    }
+    deletion.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+  }
+
+  @Test(dataProvider = "segmentTypes")
+  public void testSegmentCallbackDoesNotHoldLifecycleLock(boolean consuming)
+      throws Exception {
+    TableDataManager manager = mock(TableDataManager.class);
+    CountDownLatch addEntered = new CountDownLatch(1);
+    CountDownLatch finishAdd = new CountDownLatch(1);
+    if (consuming) {
+      doAnswer(invocation -> {
+        addEntered.countDown();
+        await(finishAdd);
+        return null;
+      }).when(manager).addConsumingSegment(anyString());
+    } else {
+      doAnswer(invocation -> {
+        addEntered.countDown();
+        await(finishAdd);
+        return null;
+      }).when(manager).addOnlineSegment(anyString());
+    }
+    TestInstanceDataManager instance = new TestInstanceDataManager(name -> 
manager);
+    FutureTask<Void> add = new FutureTask<>(() -> {
+      addSegment(instance, "callback_REALTIME", "segment", consuming);
+      return null;
+    });
+    Thread addThread = new Thread(add, "add-segment");
+    FutureTask<Void> deletion = new FutureTask<>(() -> {
+      instance.deleteTable("callback_REALTIME", 1L);
+      return null;
+    });
+    Thread deletionThread = new Thread(deletion, "delete-during-add");
+    try {
+      addThread.start();
+      await(addEntered);
+      deletionThread.start();
+      deletion.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      verify(manager).shutDown();
+    } finally {
+      finishAdd.countDown();
+      join(addThread);
+      join(deletionThread);
+    }
+    add.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+  }
+
+  @Test
+  public void testShutdownFailureReleasesLifecycleLock()
+      throws Exception {
+    TableDataManager oldManager = mock(TableDataManager.class);
+    TableDataManager replacement = mock(TableDataManager.class);
+    RuntimeException failure = new IllegalStateException("shutdown failure");
+    doThrow(failure).when(oldManager).shutDown();
+    AtomicInteger creations = new AtomicInteger();
+    TestInstanceDataManager instance = new TestInstanceDataManager(name ->
+        creations.getAndIncrement() == 0 ? oldManager : replacement);
+    instance.addConsumingSegment("failure_REALTIME", "old");
+    assertThrows(IllegalStateException.class, () -> 
instance.deleteTable("failure_REALTIME", 1L));
+    assertNull(instance.getTableDataManager("failure_REALTIME"));
+    FutureTask<Void> creation = new FutureTask<>(() -> {
+      instance.addConsumingSegment("failure_REALTIME", "new");
+      return null;
+    });
+    Thread creationThread = new Thread(creation, 
"create-after-shutdown-failure");
+    try {
+      creationThread.start();
+      creation.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      verify(replacement).addConsumingSegment("new");
+    } finally {
+      join(creationThread);
+    }
+  }
+
+  @Test
+  public void testCreationFailureReleasesLifecycleLock()
+      throws Exception {
+    TableDataManager manager = mock(TableDataManager.class);
+    AtomicInteger creations = new AtomicInteger();
+    TestInstanceDataManager instance = new TestInstanceDataManager(name -> {
+      if (creations.getAndIncrement() == 0) {
+        throw new IllegalStateException("creation failure");
+      }
+      return manager;
+    });
+    assertThrows(IllegalStateException.class, () -> 
instance.addOnlineSegment("failure_REALTIME", "first"));
+    FutureTask<Void> creation = new FutureTask<>(() -> {
+      instance.addOnlineSegment("failure_REALTIME", "second");
+      return null;
+    });
+    Thread creationThread = new Thread(creation, "retry-creation");
+    try {
+      creationThread.start();
+      creation.get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+      verify(manager).addOnlineSegment("second");
+    } finally {
+      join(creationThread);
+    }
+  }
+
+  private static void addSegment(HelixInstanceDataManager instance, String 
table, String segment, boolean consuming)
+      throws Exception {
+    if (consuming) {
+      instance.addConsumingSegment(table, segment);
+    } else {
+      instance.addOnlineSegment(table, segment);
+    }
+  }
+
+  private static void await(CountDownLatch latch)
+      throws InterruptedException {
+    assertTrue(latch.await(TIMEOUT_SECONDS, TimeUnit.SECONDS), "Timed out 
waiting for controlled lifecycle operation");
+  }
+
+  private static void join(Thread thread)
+      throws InterruptedException {
+    thread.join(TimeUnit.SECONDS.toMillis(TIMEOUT_SECONDS));
+    assertFalse(thread.isAlive(), "Lifecycle worker did not terminate");
+  }
+
+  /// Supplies controlled table instances while exercising the real instance 
manager's map and lifecycle operations.
+  private static class TestInstanceDataManager extends 
HelixInstanceDataManager {
+    private final Function<String, TableDataManager> _factory;
+
+    TestInstanceDataManager(Function<String, TableDataManager> factory) {
+      _factory = factory;
+      _recentlyDeletedTables = CacheBuilder.newBuilder().build();
+    }
+
+    @Override
+    TableDataManager createTableDataManager(String tableNameWithType) {
+      return _factory.apply(tableNameWithType);
+    }
+  }
+}


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

Reply via email to