SteNicholas commented on code in PR #3781:
URL: https://github.com/apache/celeborn/pull/3781#discussion_r3803523370


##########
common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala:
##########
@@ -656,8 +656,10 @@ class CelebornConf(loadDefaults: Boolean) extends 
Cloneable with Logging with Se
   // //////////////////////////////////////////////////////
   //                      Master                         //
   // //////////////////////////////////////////////////////
-  def masterSlotAssignPolicy: SlotsAssignPolicy =
-    SlotsAssignPolicy.valueOf(get(MASTER_SLOT_ASSIGN_POLICY))
+  def masterSlotAssignPolicyName: String = get(MASTER_SLOT_ASSIGN_POLICY)
+
+  def masterSlotAssignPolicy: BuiltInSlotsAssignPolicy =
+    BuiltInSlotsAssignPolicy.valueOf(get(MASTER_SLOT_ASSIGN_POLICY))

Review Comment:
   Preserve the existing public API here. This changes the JVM return type of 
`masterSlotAssignPolicy` and removes the public `SlotsAssignPolicy` class, so 
integrations compiled against `celeborn-common` can fail with 
`NoSuchMethodError` or `ClassNotFoundException` after upgrading. Please keep 
the old enum/accessor (possibly deprecated) and add a separate string or 
built-in-only accessor for the SPI policy name.



##########
master/src/main/java/org/apache/celeborn/service/deploy/master/slotsalloc/SlotsAssignStrategyManager.java:
##########
@@ -0,0 +1,157 @@
+/*
+ * 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.celeborn.service.deploy.master.slotsalloc;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.Locale;
+import java.util.Map;
+import java.util.ServiceLoader;
+import java.util.TreeSet;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.celeborn.common.CelebornConf;
+import org.apache.celeborn.server.common.service.config.ConfigService;
+import org.apache.celeborn.server.common.service.config.SystemConfig;
+
+/**
+ * Discovers slot assignment strategy providers and atomically switches the 
active strategy after
+ * dynamic configuration updates.
+ */
+public final class SlotsAssignStrategyManager {
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(SlotsAssignStrategyManager.class);
+
+  private final CelebornConf staticConf;
+  private final ConfigService configService;
+  private final Map<String, SlotsAssignStrategyProvider> providersByName;
+
+  private volatile ConfiguredStrategy configuredStrategy;
+  private Map<String, String> appliedDynamicConfigs;
+
+  public SlotsAssignStrategyManager(CelebornConf staticConf, ConfigService 
configService) {
+    this.staticConf = staticConf.clone();
+    this.configService = configService;
+    this.providersByName =
+        Collections.unmodifiableMap(
+            
loadProviders(ServiceLoader.load(SlotsAssignStrategyProvider.class)));
+
+    this.appliedDynamicConfigs = currentDynamicConfigs();
+    this.configuredStrategy = createStrategy(appliedDynamicConfigs);
+
+    if (configService != null) {
+      configService.registerListenerOnConfigUpdate(this::reload);
+      // Close the gap between taking the initial snapshot and registering the 
listener.
+      reload();
+    }
+  }
+
+  public SlotsAssignStrategy getStrategy() {
+    return configuredStrategy.strategy;
+  }
+
+  static Map<String, SlotsAssignStrategyProvider> loadProviders(
+      Iterable<SlotsAssignStrategyProvider> providers) {
+    Map<String, SlotsAssignStrategyProvider> providersByName = new 
LinkedHashMap<>();
+    for (SlotsAssignStrategyProvider provider : providers) {
+      String providerName = provider.getName();
+      if (providerName == null || providerName.isEmpty()) {
+        throw new IllegalStateException(
+            "Slots assignment strategy provider "
+                + provider.getClass().getName()
+                + " has an empty name");
+      }
+
+      String normalizedName = providerName.toUpperCase(Locale.ROOT);
+      SlotsAssignStrategyProvider previous = 
providersByName.put(normalizedName, provider);
+      if (previous != null) {
+        throw new IllegalStateException(
+            "Multiple slots assignment strategy providers are registered for '"
+                + providerName
+                + "': "
+                + previous.getClass().getName()
+                + " and "
+                + provider.getClass().getName());
+      }
+    }
+    return providersByName;
+  }
+
+  private synchronized void reload() {
+    Map<String, String> latestDynamicConfigs = currentDynamicConfigs();
+    if (latestDynamicConfigs.equals(appliedDynamicConfigs)) {
+      return;
+    }
+
+    try {
+      ConfiguredStrategy updatedStrategy = 
createStrategy(latestDynamicConfigs);
+      configuredStrategy = updatedStrategy;
+      appliedDynamicConfigs = latestDynamicConfigs;
+      LOG.info(
+          "Reloaded slots assignment strategy provider {} after dynamic 
configuration update",
+          updatedStrategy.providerName);
+    } catch (RuntimeException e) {
+      LOG.error(
+          "Failed to reload slots assignment strategy; keeping provider {}",
+          configuredStrategy.providerName,
+          e);
+    }
+  }
+
+  private ConfiguredStrategy createStrategy(Map<String, String> 
dynamicConfigs) {
+    CelebornConf effectiveConf = staticConf.clone();
+    dynamicConfigs.forEach(effectiveConf::set);
+
+    String providerName = effectiveConf.masterSlotAssignPolicyName();
+    SlotsAssignStrategyProvider provider =
+        providersByName.get(providerName.toUpperCase(Locale.ROOT));
+    if (provider == null) {
+      throw new IllegalArgumentException(
+          "No slots assignment strategy provider is registered for '"
+              + providerName
+              + "'. Available providers: "
+              + new TreeSet<>(providersByName.keySet()));
+    }
+    return new ConfiguredStrategy(providerName, 
provider.create(effectiveConf));

Review Comment:
   Validate that `provider.create(effectiveConf)` is non-null before publishing 
this `ConfiguredStrategy`. A third-party provider returning null is currently 
treated as a successful reload; `getStrategy()` then returns null and every 
subsequent slot allocation fails. Using `Objects.requireNonNull` here lets the 
existing reload exception path retain the last valid strategy.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to