This is an automated email from the ASF dual-hosted git repository.
morrySnow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 996cb0c37de [fix](fe) Gate serving until startup initialization
completes (#68330)
996cb0c37de is described below
commit 996cb0c37deeeb8f73aa885bd825faa6c77187aa
Author: morrySnow <[email protected]>
AuthorDate: Thu Sep 24 18:15:56 2026 +0800
[fix](fe) Gate serving until startup initialization completes (#68330)
### What problem does this PR solve?
Related PR: #68302
Problem Summary:
If UNKNOWN interrupts a FE's first FOLLOWER/OBSERVER transition,
initialization returns early while the previous FE type is retained. The
replayer can subsequently set the metadata readiness/readability flags
to true before initialization has completed. This can release
`waitForReady()` and allow local reads while startup work is still
incomplete. Clearing those flags once does not close the race, since the
replayer can set them again.
Add a process-local `startupInitialized` gate to `isReady()` and
`canRead()`. The state listener opens the gate only after a successful
MASTER/FOLLOWER/OBSERVER initialization and publication of the new FE
type. UNKNOWN handling cannot open it. Initialization's own waits
continue to use metadata readiness, avoiding a circular dependency.
Once initialization has completed, the gate stays open so an initialized
FE entering UNKNOWN retains its existing read policy and metadata
freshness checks.
### Release note
Fix a race that could let an FE report ready or serve local reads before
startup initialization completed after an UNKNOWN state interruption.
---
.../main/java/org/apache/doris/catalog/Env.java | 42 ++++--
.../apache/doris/catalog/EnvStateListenerTest.java | 167 ++++++++++++++++++++-
2 files changed, 196 insertions(+), 13 deletions(-)
diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
index ab594e7487d..15557da09bc 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/Env.java
@@ -433,12 +433,14 @@ public class Env {
protected boolean isFirstTimeStartUp = false;
protected boolean isElectable;
- // set to true after finished replay all meta and ready to serve
- // set to false when catalog is not ready.
+ // Metadata readiness, updated by the replayer independently of startup
initialization.
private AtomicBoolean isReady = new AtomicBoolean(false);
+ // Published after the first successful MASTER/FOLLOWER/OBSERVER
initialization and FE type commit.
+ // Keep this true across UNKNOWN transitions so initialized nodes retain
their existing read policy.
+ private volatile boolean startupInitialized = false;
// set to true after http server start
private AtomicBoolean httpReady = new AtomicBoolean(false);
- // set to true if FE can offer READ service.
+ // Metadata read eligibility; serving reads also requires
startupInitialized.
// canRead can be true even if isReady is false.
// for example: OBSERVER transfer to UNKNOWN, then isReady will be set to
false, but canRead can still be true
private AtomicBoolean canRead = new AtomicBoolean(false);
@@ -1308,13 +1310,18 @@ public class Env {
Thread.sleep(100);
if (counter++ % 100 == 0) {
String reason = editLog == null ? "editlog is null" :
editLog.getNotReadyReason();
- LOG.info("wait catalog to be ready. feType:{} isReady:{},
counter:{} reason: {}",
- feType, isReady.get(), counter, reason);
+ LOG.info("wait catalog to be ready. feType:{} metadataReady:{}
startupInitialized:{}, "
+ + "counter:{} reason: {}",
+ feType, isMetadataReady(), startupInitialized,
counter, reason);
}
}
}
public boolean isReady() {
+ return startupInitialized && isMetadataReady();
+ }
+
+ private boolean isMetadataReady() {
return isReady.get();
}
@@ -1957,7 +1964,8 @@ public class Env {
*/
public boolean postProcessAfterMetadataReplayed(boolean waitCatalogReady) {
if (waitCatalogReady) {
- while (!isReady()) {
+ // Startup initialization itself must not wait for the serving
gate that it will open.
+ while (!isMetadataReady()) {
// Avoid endless waiting if the state has changed.
//
// Consider the following situation:
@@ -2138,7 +2146,7 @@ public class Env {
replayer.start();
}
- // 'isReady' will be set to true in 'setCanRead()' method
+ // The replayer publishes metadata readiness before startup
initialization completes.
if (!postProcessAfterMetadataReplayed(true)) {
// A newer BDB state is already waiting in typeTransferQueue.
Abort this stale transition so the
// state listener can process the newer state instead of
waiting indefinitely for this node to
@@ -2214,7 +2222,7 @@ public class Env {
// After the cluster initialization is complete, 'lower_case_table_names'
can not be modified during the cluster
// restart or upgrade.
private void checkLowerCaseTableNames() {
- while (!isReady()) {
+ while (!isMetadataReady()) {
// Waiting for lower_case_table_names to initialize value from
image or editlog.
try {
LOG.info("Waiting for \'lower_case_table_names\'
initialization.");
@@ -3268,7 +3276,12 @@ public class Env {
}
public void startStateListener() {
- listener = new Daemon("stateListener", STATE_CHANGE_CHECK_INTERVAL_MS)
{
+ listener = createStateListener();
+ listener.start();
+ }
+
+ Daemon createStateListener() {
+ Daemon stateListener = new Daemon("stateListener",
STATE_CHANGE_CHECK_INTERVAL_MS) {
@Override
protected synchronized void runOneCycle() {
@@ -3382,13 +3395,18 @@ public class Env {
continue;
}
feType = newType;
+ // INIT -> UNKNOWN is a completed no-op, not a completed
startup initialization.
+ if (newType == FrontendNodeType.MASTER || newType ==
FrontendNodeType.FOLLOWER
+ || newType == FrontendNodeType.OBSERVER) {
+ startupInitialized = true;
+ }
LOG.info("finished to transfer FE type to {}", feType);
}
} // end runOneCycle
};
- listener.setMetaContext(metaContext);
- listener.start();
+ stateListener.setMetaContext(metaContext);
+ return stateListener;
}
public synchronized boolean replayJournal(long toJournalId) {
@@ -5539,7 +5557,7 @@ public class Env {
}
public boolean canRead() {
- return this.canRead.get();
+ return startupInitialized && canRead.get();
}
public boolean isElectable() {
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/catalog/EnvStateListenerTest.java
b/fe/fe-core/src/test/java/org/apache/doris/catalog/EnvStateListenerTest.java
index 0a41817f74b..336b6fa652a 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/catalog/EnvStateListenerTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/catalog/EnvStateListenerTest.java
@@ -17,22 +17,49 @@
package org.apache.doris.catalog;
+import org.apache.doris.common.Config;
import org.apache.doris.common.util.Daemon;
+import org.apache.doris.datasource.CatalogMgr;
import org.apache.doris.ha.FrontendNodeType;
+import org.apache.doris.metric.MetricRepo;
+import org.apache.doris.mysql.privilege.Auth;
+import org.apache.doris.statistics.analysis.AnalysisManager;
+import org.apache.doris.statistics.analysis.FollowerColumnSender;
+import org.apache.doris.statistics.cache.StatisticsCache;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.mockito.MockedConstruction;
+import org.mockito.MockedStatic;
import org.mockito.Mockito;
+import org.mockito.stubbing.Answer;
import java.lang.reflect.Field;
+import java.lang.reflect.Method;
+import java.util.ArrayList;
+import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+@Timeout(30)
public class EnvStateListenerTest {
+ private Env env;
+
+ @BeforeEach
+ public void setUp() {
+ // Cold Env initialization and Mockito instrumentation can exceed the
test timeout in CI.
+ // Keep fixture creation outside the timeout that guards the state
listener transitions.
+ env = Mockito.spy(new Env(false));
+ }
+
@Test
public void testInterruptedNonMasterTransitionDoesNotCommitFeType() throws
Exception {
- Env env = Mockito.spy(new Env(false));
setField(env, "replayer", Mockito.mock(Daemon.class));
CountDownLatch firstTransitionInterrupted = new CountDownLatch(1);
@@ -63,9 +90,147 @@ public class EnvStateListenerTest {
Assertions.assertEquals(FrontendNodeType.INIT, env.getFeType());
} finally {
stateListener.exit();
+ env.notifyNewFETypeTransfer(env.getFeType());
+ stateListener.join(5000);
+ Assertions.assertFalse(stateListener.isAlive());
+ }
+ }
+
+ @ParameterizedTest
+ @CsvSource({"INIT, FOLLOWER", "INIT, OBSERVER", "UNKNOWN, FOLLOWER",
"UNKNOWN, OBSERVER"})
+ public void
testUnknownInterruptionKeepsStartupGateClosedUntilRetryCompletes(
+ FrontendNodeType initialType, FrontendNodeType targetType) throws
Exception {
+ setField(env, "feType", initialType);
+ Mockito.doReturn(false).when(env).replayJournal(-1);
+ env.createReplayer();
+ Daemon replayer = (Daemon) getField(env, "replayer");
+ Daemon stateListener = env.createStateListener();
+
+ Auth auth = Mockito.mock(Auth.class);
+ setField(env, "auth", auth);
+ CatalogMgr catalogMgr = Mockito.spy(env.getCatalogMgr());
+ setField(env, "catalogMgr", catalogMgr);
+ AnalysisManager analysisManager = Mockito.mock(AnalysisManager.class);
+ StatisticsCache statisticsCache = Mockito.mock(StatisticsCache.class);
+
Mockito.when(analysisManager.getStatisticsCache()).thenReturn(statisticsCache);
+ setField(env, "analysisManager", analysisManager);
+
+ List<Boolean> servingDuringInitialization = new ArrayList<>();
+ Answer<Void> recordServingState = invocation -> {
+ servingDuringInitialization.add(env.isReady() || env.canRead());
+ return null;
+ };
+
Mockito.doAnswer(recordServingState).when(env).startNonMasterDaemonThreads();
+ Mockito.doAnswer(recordServingState).when(statisticsCache).preHeat();
+
+ AtomicInteger transitionAttempts = new AtomicInteger();
+ Mockito.doAnswer(invocation -> {
+ if (transitionAttempts.incrementAndGet() == 1) {
+ boolean completed = (boolean) invocation.callRealMethod();
+ // Run the real replayer readiness update after UNKNOWN
interrupts the wait, before
+ // the listener handles UNKNOWN. This deterministically
reproduces the unsafe interleaving.
+ replayFreshMetadata(env, replayer);
+ return completed;
+ }
+ // transferToNonMaster resets metadata readiness on every attempt.
Let the replayer
+ // catch up again before the retry executes the real metadata
post-processing.
+ replayFreshMetadata(env, replayer);
+ return invocation.callRealMethod();
+ }).when(env).postProcessAfterMetadataReplayed(true);
+
+ try (MockedStatic<MetricRepo> metrics =
Mockito.mockStatic(MetricRepo.class);
+ MockedConstruction<FollowerColumnSender> senders =
Mockito.mockConstruction(
+ FollowerColumnSender.class,
+ (sender, context) ->
Mockito.doAnswer(recordServingState).when(sender).start())) {
+ metrics.when(MetricRepo::init).thenAnswer(recordServingState);
+
+ env.notifyNewFETypeTransfer(targetType);
+ env.notifyNewFETypeTransfer(FrontendNodeType.UNKNOWN);
+ // An equal event ends runOneCycle. For INIT, first commit
UNKNOWN; for an initial
+ // UNKNOWN, the interruption event itself is already equal to the
retained FE type.
+ if (initialType == FrontendNodeType.INIT) {
+ env.notifyNewFETypeTransfer(FrontendNodeType.UNKNOWN);
+ }
+ runOneCycle(stateListener);
+
+ Assertions.assertEquals(FrontendNodeType.UNKNOWN, env.getFeType());
+ Assertions.assertEquals(1, transitionAttempts.get());
+ Assertions.assertTrue(((AtomicBoolean) getField(env,
"isReady")).get());
+ Assertions.assertTrue(((AtomicBoolean) getField(env,
"canRead")).get());
+ Assertions.assertFalse(env.isReady());
+ Assertions.assertFalse(env.canRead());
+ Mockito.verify(auth, Mockito.never()).rectifyPrivs();
+ Mockito.verify(catalogMgr,
Mockito.never()).registerCatalogRefreshListener(env);
+ Mockito.verify(env, Mockito.never()).startNonMasterDaemonThreads();
+ Mockito.verify(statisticsCache, Mockito.never()).preHeat();
+ metrics.verify(MetricRepo::init, Mockito.never());
+ Assertions.assertTrue(senders.constructed().isEmpty());
+
+ replayFreshMetadata(env, replayer);
+ Assertions.assertFalse(env.isReady());
+ Assertions.assertFalse(env.canRead());
+
+ env.notifyNewFETypeTransfer(targetType);
+ env.notifyNewFETypeTransfer(targetType);
+ runOneCycle(stateListener);
+
+ Assertions.assertEquals(targetType, env.getFeType());
+ Assertions.assertEquals(2, transitionAttempts.get());
+ Mockito.verify(auth).rectifyPrivs();
+ Mockito.verify(catalogMgr).registerCatalogRefreshListener(env);
+ Mockito.verify(env).startNonMasterDaemonThreads();
+ metrics.verify(MetricRepo::init);
+ Mockito.verify(statisticsCache).preHeat();
+ Assertions.assertEquals(1, senders.constructed().size());
+ Mockito.verify(senders.constructed().get(0)).start();
+ Assertions.assertEquals(List.of(false, false, false, false),
servingDuringInitialization);
+ Assertions.assertTrue(env.isReady());
+ Assertions.assertTrue(env.canRead());
+ env.waitForReady();
+
+ // Once startup has completed, UNKNOWN retains the existing read
policy until metadata expires.
+ env.notifyNewFETypeTransfer(FrontendNodeType.UNKNOWN);
+ env.notifyNewFETypeTransfer(FrontendNodeType.UNKNOWN);
+ runOneCycle(stateListener);
+ Assertions.assertFalse(env.isReady());
+ Assertions.assertTrue(env.canRead());
+ replayFreshMetadata(env, replayer);
+ Assertions.assertTrue(env.isReady());
+ Assertions.assertTrue(env.canRead());
+ env.setSynchronizedTime(0);
+ runOneCycle(replayer);
+ Assertions.assertFalse(env.isReady());
+ Assertions.assertFalse(env.canRead());
+ }
+ }
+
+ @Test
+ public void testIgnoreMetaCheckDoesNotBypassStartupGate() throws Exception
{
+ Mockito.doReturn(false).when(env).replayJournal(-1);
+ env.createReplayer();
+ boolean originalIgnoreMetaCheck = Config.ignore_meta_check;
+ try {
+ Config.ignore_meta_check = true;
+ runOneCycle((Daemon) getField(env, "replayer"));
+ Assertions.assertTrue(((AtomicBoolean) getField(env,
"canRead")).get());
+ Assertions.assertFalse(env.isReady());
+ Assertions.assertFalse(env.canRead());
+ } finally {
+ Config.ignore_meta_check = originalIgnoreMetaCheck;
}
}
+ private static void replayFreshMetadata(Env env, Daemon replayer) throws
ReflectiveOperationException {
+ env.setSynchronizedTime(System.currentTimeMillis());
+ runOneCycle(replayer);
+ }
+
+ private static void runOneCycle(Daemon daemon) throws
ReflectiveOperationException {
+ Method method = daemon.getClass().getDeclaredMethod("runOneCycle");
+ method.setAccessible(true);
+ method.invoke(daemon);
+ }
+
private static Object getField(Env env, String fieldName) throws
ReflectiveOperationException {
Field field = Env.class.getDeclaredField(fieldName);
field.setAccessible(true);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]