rzo1 commented on PR #9153:
URL: https://github.com/apache/storm/pull/9153#issuecomment-6057106209

   Thanks for the log @jkrauss82, that's helpful.
   
   I'd rather not drop the `throw` for the whole `run()`. That would also 
swallow persistent local failures (e.g. `mkSlot`, or the "multiple topologies 
assigned to one port" check), and `readAssignments()` shows that failing fast 
after a few retries is intended there.
   
   Looking at the trace, though, the supervisor does not read assignments from 
ZooKeeper at this point: `assignments()` / `assignmentsInfo()` come from the 
in-memory local backend, which is filled from Nimbus via Thrift. The only 
ZooKeeper read in `ReadClusterState.run()` is the profiler request lookup 
(`getProfileActions` → `getTopologyProfileRequests` → `node_exists`). So right 
now a ZK hiccup on an optional feature takes down the supervisor and all its 
workers.
   
   I'd suggest scoping the `catch` to that call instead. Skipping a cycle is 
safe: `Slot.addProfilerActions` only adds actions, and the requests stay in 
ZooKeeper, so the next sync (every 10s) picks them up.
   
   ```diff
   @@ ReadClusterState.run() @@
                    //Something odd happened try again later
                    return;
                }
   -            Map<String, List<ProfileRequest>> topoIdToProfilerActions = 
getProfileActions(stormClusterState, stormIds);
   +            Map<String, List<ProfileRequest>> topoIdToProfilerActions;
   +            try {
   +                topoIdToProfilerActions = 
getProfileActions(stormClusterState, stormIds);
   +            } catch (Exception e) {
   +                // Assignments come from the local backend; profiler 
requests are the only ZooKeeper read here.
   +                // They stay in ZooKeeper, so skip them for this cycle 
instead of halting the supervisor.
   +                LOG.warn("Failed to read profiler requests, will retry next 
cycle", e);
   +                topoIdToProfilerActions = Collections.emptyMap();
   +            }
    
                HashSet<Integer> assignedPorts = new HashSet<>();
                LOG.debug("Synchronizing supervisor");
   ```
   (plus `import java.util.Collections;`)
   
   Here's a test for it. Without the change, 
`profilerRequestFailureDoesNotHaltAssignmentSync` fails with the same 
`RuntimeException` → `SessionExpiredException` chain as in your log. With it, 
both tests pass.
   
   
<details><summary><code>storm-server/src/test/java/org/apache/storm/daemon/supervisor/ReadClusterStateTest.java</code></summary>
   
   ```java
   /*
    * 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.storm.daemon.supervisor;
   
   import java.nio.file.Path;
   import java.util.Collections;
   import java.util.HashMap;
   import java.util.Map;
   import java.util.concurrent.atomic.AtomicReference;
   import org.apache.storm.Config;
   import org.apache.storm.DaemonConfig;
   import org.apache.storm.cluster.IStormClusterState;
   import org.apache.storm.scheduler.ISupervisor;
   import org.apache.storm.shade.org.apache.zookeeper.KeeperException;
   import org.junit.jupiter.api.Test;
   import org.junit.jupiter.api.io.TempDir;
   
   import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
   import static org.junit.jupiter.api.Assertions.assertThrows;
   import static org.mockito.ArgumentMatchers.any;
   import static org.mockito.ArgumentMatchers.anyString;
   import static org.mockito.Mockito.mock;
   import static org.mockito.Mockito.never;
   import static org.mockito.Mockito.verify;
   import static org.mockito.Mockito.when;
   
   public class ReadClusterStateTest {
   
       private static final String TOPOLOGY_ID = "topo-1";
   
       @TempDir
       Path localDir;
   
       private final IStormClusterState clusterState = 
mock(IStormClusterState.class);
       private final ISupervisor iSupervisor = mock(ISupervisor.class);
   
       private ReadClusterState newReadClusterState() throws Exception {
           Map<String, Object> conf = new HashMap<>();
           conf.put(Config.STORM_CLUSTER_MODE, "local");
           conf.put(Config.STORM_LOCAL_DIR, localDir.toString());
           conf.put(DaemonConfig.SUPERVISOR_SLOTS_PORTS, 
Collections.emptyList());
   
           Supervisor supervisor = mock(Supervisor.class);
           when(supervisor.getConf()).thenReturn(conf);
           when(supervisor.getStormClusterState()).thenReturn(clusterState);
           when(supervisor.getAssignmentId()).thenReturn("supervisor-1");
           when(supervisor.getiSupervisor()).thenReturn(iSupervisor);
           when(supervisor.getHostName()).thenReturn("host-1");
           when(supervisor.getCurrAssignment()).thenReturn(new 
AtomicReference<>(new HashMap<>()));
           return new ReadClusterState(supervisor);
       }
   
       @Test
       public void profilerRequestFailureDoesNotHaltAssignmentSync() throws 
Exception {
           
when(clusterState.assignments(null)).thenReturn(Collections.singletonList(TOPOLOGY_ID));
           
when(clusterState.assignmentsInfo()).thenReturn(Collections.emptyMap());
           // What ClientZookeeper surfaces after Curator gives up, as in the 
reported supervisor crash.
           
when(clusterState.getTopologyProfileRequests(TOPOLOGY_ID)).thenThrow(new 
RuntimeException(
               new KeeperException.SessionExpiredException()));
   
           ReadClusterState readClusterState = newReadClusterState();
   
           assertDoesNotThrow(readClusterState::run);
           verify(iSupervisor).assigned(any());
       }
   
       @Test
       public void assignmentReadFailureStillPropagates() throws Exception {
           when(clusterState.assignments(null)).thenThrow(new 
IllegalStateException("broken local backend"));
   
           ReadClusterState readClusterState = newReadClusterState();
   
           assertThrows(RuntimeException.class, readClusterState::run);
           verify(clusterState, 
never()).getTopologyProfileRequests(anyString());
       }
   }
   ```
   </details>
   
   Side note: this log supports @GGraziadei's point. 
`ConnectionAwareRetryPolicy` held the heartbeat thread for 20s and then gave up 
anyway. ZooKeeper was unavailable for longer than the session timeout, so 
waiting couldn't help, and it was the `catch` in `SupervisorHeartbeat` that 
kept the process alive. I'd keep this PR to the two scoped `catch` blocks 
(heartbeat + profiler lookup) and move `ConnectionAwareRetryPolicy` to a 
follow-up with @reiabreu's foreground/background fix and tests against a real 
ZooKeeper failover.
   


-- 
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