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]

Reply via email to