Repository: nifi
Updated Branches:
  refs/heads/master 4544f3969 -> 8acac9cba


NIFI-5153: If a node is disconnected due to failure to complete mutable 
request, the node should be allowed to rejoin

This closes #2677.

Signed-off-by: Bryan Bende <[email protected]>


Project: http://git-wip-us.apache.org/repos/asf/nifi/repo
Commit: http://git-wip-us.apache.org/repos/asf/nifi/commit/8acac9cb
Tree: http://git-wip-us.apache.org/repos/asf/nifi/tree/8acac9cb
Diff: http://git-wip-us.apache.org/repos/asf/nifi/diff/8acac9cb

Branch: refs/heads/master
Commit: 8acac9cba5b9d409e4ba885ce5dba15ef416de0e
Parents: 4544f39
Author: Mark Payne <[email protected]>
Authored: Fri May 4 09:31:51 2018 -0400
Committer: Bryan Bende <[email protected]>
Committed: Mon May 7 15:41:38 2018 -0400

----------------------------------------------------------------------
 .../cluster/coordination/node/NodeClusterCoordinator.java   | 9 ++++++---
 1 file changed, 6 insertions(+), 3 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/nifi/blob/8acac9cb/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/node/NodeClusterCoordinator.java
----------------------------------------------------------------------
diff --git 
a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/node/NodeClusterCoordinator.java
 
b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/node/NodeClusterCoordinator.java
index 754d370..4e4625c 100644
--- 
a/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/node/NodeClusterCoordinator.java
+++ 
b/nifi-nar-bundles/nifi-framework-bundle/nifi-framework/nifi-framework-cluster/src/main/java/org/apache/nifi/cluster/coordination/node/NodeClusterCoordinator.java
@@ -373,7 +373,7 @@ public class NodeClusterCoordinator implements 
ClusterCoordinator, ProtocolHandl
         final Map<NodeConnectionState, List<NodeIdentifier>> connectionStates 
= new HashMap<>();
         for (final Map.Entry<NodeIdentifier, NodeConnectionStatus> entry : 
nodeStatuses.entrySet()) {
             final NodeConnectionState state = entry.getValue().getState();
-            final List<NodeIdentifier> nodeIds = 
connectionStates.computeIfAbsent(state, s -> new ArrayList<NodeIdentifier>());
+            final List<NodeIdentifier> nodeIds = 
connectionStates.computeIfAbsent(state, s -> new ArrayList<>());
             nodeIds.add(entry.getKey());
         }
 
@@ -998,9 +998,12 @@ public class NodeClusterCoordinator implements 
ClusterCoordinator, ProtocolHandl
             // disconnect problematic nodes
             if (!problematicNodeResponses.isEmpty() && 
problematicNodeResponses.size() < nodeResponses.size()) {
                 final Set<NodeIdentifier> failedNodeIds = 
problematicNodeResponses.stream().map(response -> 
response.getNodeId()).collect(Collectors.toSet());
-                logger.warn(String.format("The following nodes failed to 
process URI %s '%s'.  Requesting each node disconnect from cluster.", uriPath, 
failedNodeIds));
+                logger.warn(String.format("The following nodes failed to 
process URI %s '%s'.  Requesting each node reconnect to cluster.", uriPath, 
failedNodeIds));
                 for (final NodeIdentifier nodeId : failedNodeIds) {
-                    requestNodeDisconnect(nodeId, 
DisconnectionCode.FAILED_TO_SERVICE_REQUEST, "Failed to process request " + 
method + " " + uriPath);
+                    // Update the node to 'CONNECTING' status and request that 
the node connect
+                    final NodeConnectionStatus reconnectionStatus = new 
NodeConnectionStatus(nodeId, NodeConnectionState.CONNECTING);
+                    updateNodeStatus(reconnectionStatus);
+                    requestNodeConnect(nodeId, null);
                 }
             }
         }

Reply via email to