This is an automated email from the ASF dual-hosted git repository.

pgyori pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/nifi.git


The following commit(s) were added to refs/heads/main by this push:
     new ceb2838f161 NIFI-14737: Add Flow Analysis Rule to Require 
Load-Balanced Connections from a Primary-Only Source Processor (#10081)
ceb2838f161 is described below

commit ceb2838f1610e7ab5d34c546f90835b08ebb94a0
Author: Matt Burgess <[email protected]>
AuthorDate: Mon Aug 3 10:49:30 2026 -0400

    NIFI-14737: Add Flow Analysis Rule to Require Load-Balanced Connections 
from a Primary-Only Source Processor (#10081)
    
    * NIFI-14737: Add Flow Analysis Rule to Require Load-Balanced Connections 
from a Source Processor running only on the Primary Node
    
    * Incorporated review comments
---
 ...dConnectionAfterPrimaryNodeSourceProcessor.java | 77 ++++++++++++++++++++++
 .../org.apache.nifi.flowanalysis.FlowAnalysisRule  |  1 +
 ...nectionAfterPrimaryNodeSourceProcessorTest.java | 54 +++++++++++++++
 ...fterNonPrimarySourceProcessor_No_Violation.json |  1 +
 ...onnectionAfterSourceProcessor_No_Violation.json |  1 +
 ...edConnectionAfterSourceProcessor_Violation.json |  1 +
 6 files changed, 135 insertions(+)

diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor.java
new file mode 100644
index 00000000000..82cdb676f2e
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor.java
@@ -0,0 +1,77 @@
+/*
+ *  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.nifi.flowanalysis.rules;
+
+import org.apache.nifi.annotation.documentation.CapabilityDescription;
+import org.apache.nifi.annotation.documentation.Tags;
+import org.apache.nifi.flow.ConnectableComponentType;
+import org.apache.nifi.flow.VersionedProcessGroup;
+import org.apache.nifi.flow.VersionedProcessor;
+import org.apache.nifi.flowanalysis.AbstractFlowAnalysisRule;
+import org.apache.nifi.flowanalysis.FlowAnalysisRuleContext;
+import org.apache.nifi.flowanalysis.GroupAnalysisResult;
+import org.apache.nifi.scheduling.ExecutionNode;
+
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+
+@Tags({"connection", "source", "load", "balance", "primary"})
+@CapabilityDescription("Produces rule violation when a source processor 
running only on the Primary Node does not have a Load-Balanced Connection (LBC) 
downstream"
+        + " in all outgoing connections")
+public class RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor 
extends AbstractFlowAnalysisRule {
+
+
+    @Override
+    public Collection<GroupAnalysisResult> analyzeProcessGroup(final 
VersionedProcessGroup processGroup, final FlowAnalysisRuleContext context) {
+
+        final Collection<GroupAnalysisResult> analysisResults = new 
HashSet<>();
+
+        Map<String, VersionedProcessor> sourceProcessors = new HashMap<>();
+        for (VersionedProcessor processor : processGroup.getProcessors()) {
+            if 
(ExecutionNode.PRIMARY.name().equals(processor.getExecutionNode())) {
+                sourceProcessors.put(processor.getIdentifier(), processor);
+            }
+        }
+        // Remove all processors that are not source processors or not running
+        processGroup.getConnections().forEach(connection -> {
+            if (connection.getDestination() != null && 
connection.getDestination().getType() == ConnectableComponentType.PROCESSOR) {
+                sourceProcessors.remove(connection.getDestination().getId());
+            }
+        });
+
+        processGroup.getConnections().forEach(connection -> {
+            if (connection.getSource() != null && 
connection.getSource().getId() != null) {
+                VersionedProcessor sourceProcessor = 
sourceProcessors.get(connection.getSource().getId());
+                if (sourceProcessor != null) {
+                    // Check if the connection is a Load-Balanced Connection
+                    if (connection.getLoadBalanceStrategy() == null || 
"DO_NOT_LOAD_BALANCE".equals(connection.getLoadBalanceStrategy())) {
+                        // If not, we need to report this as a violation
+                        analysisResults.add(
+                                
GroupAnalysisResult.forComponent(sourceProcessor,
+                                                
String.format("%s_LoadBalancedConnectionRequired", connection.getIdentifier()),
+                                                "Source processors running on 
the Primary Node only must configure their downstream connections to use a Load 
Balancing Strategy")
+                                        .build());
+                    }
+                }
+            }
+        });
+
+        return analysisResults;
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/resources/META-INF/services/org.apache.nifi.flowanalysis.FlowAnalysisRule
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/resources/META-INF/services/org.apache.nifi.flowanalysis.FlowAnalysisRule
index 165b44830c1..39371798eb4 100644
--- 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/resources/META-INF/services/org.apache.nifi.flowanalysis.FlowAnalysisRule
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/main/resources/META-INF/services/org.apache.nifi.flowanalysis.FlowAnalysisRule
@@ -14,6 +14,7 @@
 # limitations under the License.
 
 org.apache.nifi.flowanalysis.rules.DisallowComponentType
+org.apache.nifi.flowanalysis.rules.RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor
 org.apache.nifi.flowanalysis.rules.RequireServerSSLContextService
 org.apache.nifi.flowanalysis.rules.RestrictBackpressureSettings
 org.apache.nifi.flowanalysis.rules.RestrictFlowFileExpiration
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessorTest.java
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessorTest.java
new file mode 100644
index 00000000000..b5972dc2437
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/java/org/apache/nifi/flowanalysis/rules/RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessorTest.java
@@ -0,0 +1,54 @@
+/*
+ * 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.nifi.flowanalysis.rules;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.List;
+
+
+public class RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessorTest 
extends 
AbstractFlowAnalysisRuleTest<RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor>
 {
+
+    @Test
+    public void testViolations() throws Exception {
+        testAnalyzeProcessGroup(
+                
"src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_Violation.json",
+                List.of(
+                        "f54fbff5-0197-1000-1c04-924d2ac6cc8e" // processor 
GenerateFlowFile connecting to a funnel with no load-balancing strategy
+                )
+        );
+    }
+
+    @Test
+    public void testNoViolations() throws Exception {
+        testAnalyzeProcessGroup(
+                
"src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_No_Violation.json",
+                Collections.emptyList()
+        );
+
+        testAnalyzeProcessGroup(
+                
"src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterNonPrimarySourceProcessor_No_Violation.json",
+                Collections.emptyList()
+        );
+    }
+
+    @Override
+    protected RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor 
initializeRule() {
+        return new 
RequireLoadBalancedConnectionAfterPrimaryNodeSourceProcessor();
+    }
+}
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterNonPrimarySourceProcessor_No_Violation.json
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterNonPrimarySourceProcessor_No_Violation.json
new file mode 100644
index 00000000000..ec7ece95d4b
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterNonPrimarySourceProcessor_No_Violation.json
@@ -0,0 +1 @@
+{"flowContents":{"identifier":"29d2c039-8aca-3dde-97cc-3e13c02610ef","instanceIdentifier":"f54e4e9a-0197-1000-5827-898668a4b681","name":"NiFi
 
Flow","comments":"","position":{"x":0.0,"y":0.0},"processGroups":[],"remoteProcessGroups":[],"processors":[{"identifier":"95fc919d-97cd-3401-8b55-353969af445b","instanceIdentifier":"f54fbff5-0197-1000-1c04-924d2ac6cc8e","name":"GenerateFlowFile","comments":"","position":{"x":8.0,"y":-248.0},"type":"org.apache.nifi.processors.standard.GenerateFlowFi
 [...]
\ No newline at end of file
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_No_Violation.json
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_No_Violation.json
new file mode 100644
index 00000000000..3a5669d99e2
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_No_Violation.json
@@ -0,0 +1 @@
+{"flowContents":{"identifier":"29d2c039-8aca-3dde-97cc-3e13c02610ef","instanceIdentifier":"f54e4e9a-0197-1000-5827-898668a4b681","name":"NiFi
 
Flow","comments":"","position":{"x":0.0,"y":0.0},"processGroups":[],"remoteProcessGroups":[],"processors":[{"identifier":"95fc919d-97cd-3401-8b55-353969af445b","instanceIdentifier":"f54fbff5-0197-1000-1c04-924d2ac6cc8e","name":"GenerateFlowFile","comments":"","position":{"x":8.0,"y":-248.0},"type":"org.apache.nifi.processors.standard.GenerateFlowFi
 [...]
\ No newline at end of file
diff --git 
a/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_Violation.json
 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_Violation.json
new file mode 100644
index 00000000000..75c7344faf5
--- /dev/null
+++ 
b/nifi-extension-bundles/nifi-standard-bundle/nifi-standard-rules/src/test/resources/RequireLoadBalancedConnectionAfterSourceProcessor/RequireLoadBalancedConnectionAfterSourceProcessor_Violation.json
@@ -0,0 +1 @@
+{"flowContents":{"identifier":"29d2c039-8aca-3dde-97cc-3e13c02610ef","instanceIdentifier":"f54e4e9a-0197-1000-5827-898668a4b681","name":"NiFi
 
Flow","comments":"","position":{"x":0.0,"y":0.0},"processGroups":[],"remoteProcessGroups":[],"processors":[{"identifier":"95fc919d-97cd-3401-8b55-353969af445b","instanceIdentifier":"f54fbff5-0197-1000-1c04-924d2ac6cc8e","name":"GenerateFlowFile","comments":"","position":{"x":8.0,"y":-248.0},"type":"org.apache.nifi.processors.standard.GenerateFlowFi
 [...]
\ No newline at end of file

Reply via email to