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); } } }
