This is an automated email from the ASF dual-hosted git repository.
czy006 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/amoro.git
The following commit(s) were added to refs/heads/master by this push:
new 5d976fc71 [Refactor] Unify the main process of amoro service container
(#4339)
5d976fc71 is described below
commit 5d976fc71445690bb557e495d4787357040a1e8d
Author: ZhouJinsong <[email protected]>
AuthorDate: Mon Aug 31 17:01:31 2026 +0800
[Refactor] Unify the main process of amoro service container (#4339)
* Unify the main process of amoro service container
* Fix checkstyle issus
* Stop leader service first in dispose path
---
.../apache/amoro/server/AmoroServiceContainer.java | 148 ++++++++++++---------
.../org/apache/amoro/server/AmsEnvironment.java | 4 +-
2 files changed, 90 insertions(+), 62 deletions(-)
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
index 2e72cef20..bea5498e9 100644
--- a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
+++ b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
@@ -131,6 +131,7 @@ public class AmoroServiceContainer {
public AmoroServiceContainer() throws Exception {
initConfig();
haContainer = HighAvailabilityContainerFactory.create(serviceConfig);
+ haContainer.registerAndElect();
}
public static void main(String[] args) {
@@ -144,28 +145,21 @@ public class AmoroServiceContainer {
service.dispose();
LOG.info("AMS service has been shut down");
}));
- service.startRestServices();
- if (IS_MASTER_SLAVE_MODE) {
- // Even if one does not become the master, it cannot block the
subsequent logic.
- service.registAndElect();
- // Regardless of whether tp becomes the master, the service needs to
be activated.
- service.startOptimizingService();
- } else {
- while (true) {
- try {
- // Used to block AMS instances that have not acquired leadership
- service.waitLeaderShip();
- service.transitionToLeader();
- // Used to block AMS instances that have acquired leadership
- service.waitFollowerShip();
- } catch (ConfigurationException e) {
- LOG.error("AMS will exit...", e);
- System.exit(1);
- } catch (Exception e) {
- LOG.error("AMS start error", e);
- } finally {
- service.transitionToFollower();
- }
+
+ service.startBaseServices();
+
+ while (true) {
+ try {
+ service.waitLeaderShip();
+ service.startLeaderServices();
+ service.waitFollowerShip();
+ } catch (ConfigurationException e) {
+ LOG.error("AMS will exit...", e);
+ System.exit(1);
+ } catch (Exception e) {
+ LOG.error("AMS start error", e);
+ } finally {
+ service.stopLeaderServices();
}
}
} catch (Throwable t) {
@@ -174,10 +168,6 @@ public class AmoroServiceContainer {
}
}
- public void registAndElect() throws Exception {
- haContainer.registerAndElect();
- }
-
public enum HAState {
INITIALIZING(0),
FOLLOWER(1),
@@ -220,24 +210,25 @@ public class AmoroServiceContainer {
registerAmsServiceMetric();
}
- public void transitionToLeader() throws Exception {
- if (haState == HAState.LEADER) {
- return;
- }
- startOptimizingService();
- haState = HAState.LEADER;
- }
-
- public void transitionToFollower() {
- if (haState == HAState.FOLLOWER) {
- return;
+ /**
+ * Start base services that every AMS node needs regardless of HA mode:
REST/HTTP and
+ * catalog/table managers. In master-slave mode, also starts the optimizing
service (including
+ * Thrift), because non-leader nodes must serve optimizer requests. In
active-standby mode, the
+ * optimizing service is deferred to {@link #startLeaderServices} since the
standby node does not
+ * serve any requests.
+ */
+ public void startBaseServices() throws Exception {
+ startRestServices();
+ if (IS_MASTER_SLAVE_MODE) {
+ startOptimizingService();
}
- haState = HAState.FOLLOWER;
- disposeOptimizingService();
}
- public void startOptimizingService() throws Exception {
-
+ /**
+ * Create optimizing service objects, register handler chains, initialize
table service, and start
+ * the Thrift servers.
+ */
+ private void startOptimizingService() throws Exception {
// Load process factories and build action coordinators from default table
runtime factory.
TableProcessFactoryManager tableProcessFactoryManager = new
TableProcessFactoryManager();
tableProcessFactoryManager.initialize();
@@ -255,20 +246,6 @@ public class AmoroServiceContainer {
bucketAssignStore = BucketAssignStoreFactory.create(haContainer,
serviceConfig);
}
- // In master-slave mode, create AmsAssignService for bucket assignment
(shares BucketAssignStore
- // with DefaultTableService).
- if (IS_MASTER_SLAVE_MODE && haContainer != null && bucketAssignStore !=
null) {
- try {
- amsAssignService = new AmsAssignService(haContainer, serviceConfig,
bucketAssignStore);
- amsAssignService.start();
- LOG.info("AmsAssignService started for master-slave mode");
- } catch (UnsupportedOperationException e) {
- LOG.info("Skip AmsAssignService: {}", e.getMessage());
- } catch (Exception e) {
- LOG.error("Failed to start AmsAssignService", e);
- }
- }
-
List<ActionCoordinator> actionCoordinators =
defaultRuntimeFactory.supportedCoordinators();
tableService =
@@ -300,6 +277,61 @@ public class AmoroServiceContainer {
startThriftService();
}
+ /**
+ * Start leader-exclusive services. In active-standby mode this starts the
optimizing service
+ * (including Thrift) since only the leader serves requests. In master-slave
mode the optimizing
+ * service is already started in {@link #startBaseServices} because
non-leader nodes also serve
+ * optimizer requests; here only the leader-exclusive schedulers (e.g.
AmsAssignService) are
+ * started.
+ */
+ public void startLeaderServices() throws Exception {
+ if (haState == HAState.LEADER) {
+ return;
+ }
+ if (IS_MASTER_SLAVE_MODE) {
+ // AmsAssignService may have been stopped and set to null by a previous
stopLeaderServices
+ // call (leader re-election); recreate it if needed.
+ if (amsAssignService == null && haContainer != null) {
+ try {
+ BucketAssignStore bucketAssignStore =
+ BucketAssignStoreFactory.create(haContainer, serviceConfig);
+ amsAssignService = new AmsAssignService(haContainer, serviceConfig,
bucketAssignStore);
+ } catch (Exception e) {
+ LOG.error("Failed to recreate AmsAssignService", e);
+ }
+ }
+ if (amsAssignService != null) {
+ amsAssignService.start();
+ LOG.info("AmsAssignService started");
+ }
+ } else {
+ startOptimizingService();
+ }
+ haState = HAState.LEADER;
+ }
+
+ /**
+ * Stop leader-exclusive services. In active-standby mode this disposes the
entire optimizing
+ * service (including Thrift and service objects), since the standby node
does not serve any
+ * requests. In master-slave mode only the leader-exclusive schedulers are
stopped; Thrift stays
+ * running so the node can continue serving optimizer requests.
+ */
+ public void stopLeaderServices() {
+ if (haState == HAState.FOLLOWER) {
+ return;
+ }
+ if (IS_MASTER_SLAVE_MODE) {
+ if (amsAssignService != null) {
+ LOG.info("Stopping AmsAssignService...");
+ amsAssignService.stop();
+ amsAssignService = null;
+ }
+ } else {
+ disposeOptimizingService();
+ }
+ haState = HAState.FOLLOWER;
+ }
+
private void addHandlerChain(RuntimeHandlerChain chain) {
if (chain != null) {
tableService.addHandlerChain(chain);
@@ -316,11 +348,6 @@ public class AmoroServiceContainer {
LOG.info("Stopping optimizing server[serving:{}] ...",
optimizingServiceServer.isServing());
optimizingServiceServer.stop();
}
- if (amsAssignService != null) {
- LOG.info("Stopping AmsAssignService...");
- amsAssignService.stop();
- amsAssignService = null;
- }
if (tableService != null) {
LOG.info("Stopping table service...");
tableService.dispose();
@@ -360,6 +387,7 @@ public class AmoroServiceContainer {
}
public void dispose() {
+ stopLeaderServices();
disposeOptimizingService();
disposeRestService();
}
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/AmsEnvironment.java
b/amoro-ams/src/test/java/org/apache/amoro/server/AmsEnvironment.java
index 07898363a..db30b81b1 100644
--- a/amoro-ams/src/test/java/org/apache/amoro/server/AmsEnvironment.java
+++ b/amoro-ams/src/test/java/org/apache/amoro/server/AmsEnvironment.java
@@ -333,8 +333,8 @@ public class AmsEnvironment {
AmoroManagementConf.OPTIMIZING_SERVICE_THRIFT_BIND_PORT,
optimizingServiceBindPort);
serviceConfig.set(
AmoroManagementConf.REFRESH_EXTERNAL_CATALOGS_INTERVAL,
Duration.ofMillis(1000L));
- serviceContainer.startRestServices();
- serviceContainer.startOptimizingService();
+ serviceContainer.startBaseServices();
+ serviceContainer.startLeaderServices();
LOG.info("Started test AMS.");
break;
} catch (TTransportException e) {