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]