This is an automated email from the ASF dual-hosted git repository.
hansva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/hop.git
The following commit(s) were added to refs/heads/main by this push:
new f332853553 Issue #7067 : Kafka Consumer step fails when run inside
async web service (#7538)
f332853553 is described below
commit f332853553ad96f69a2635f31c49a2638777a290
Author: Matt Casters <[email protected]>
AuthorDate: Thu Jul 16 11:00:46 2026 +0200
Issue #7067 : Kafka Consumer step fails when run inside async web service
(#7538)
* issue #7507 : restore empty-parameter defaults after single-JVM runner fix
#7517 preferred any existing same-named variable over empty parameter
defaults and also copied parent variables into nested workflow parameters
when pass_all_parameters=Y. That broke backbone IT main-0004/main-0006
(and hop-run of those workflows alone), while empty defaults are intentional
for unset parameters.
Refine activateParameters to prefer existing variables only when the
parameter default is non-empty (HOSTNAME=localhost clobber case). Revert
ActionWorkflow parent-variable fallback; formal parent parameters still
pass through. Add unit coverage. Ignore SSH samples that need a lab host.
Verified: hop-run alone for main-0004 and main-0006 exits 0. Confirm on
Jenkins.
* issue #7519 : Kafka Consumer stop when idle
Add stopWhenIdle and maxIdleTimeMs options so the consumer can drain a
topic and exit after a configurable idle period (batch-style use cases).
* issue #7519 : fix Kafka IT bootstrap and idle assignment race
Harden BOOTSTRAP_SERVERS for the single-JVM runner, wait for Kafka
health in docker, ignore idle time until partitions are assigned, and
convert basic/mapping ITs to stop-when-idle.
* issue #7067 : Kafka Consumer step fails when run inside async web service
---
.../org/apache/hop/core/variables/Variables.java | 8 +++++
.../apache/hop/core/variables/VariablesTest.java | 19 +++++++++++
.../hop/pipeline/TransformWithMappingMeta.java | 12 +++++--
.../hop/pipeline/TransformWithMappingMetaTest.java | 39 ++++++++++++++++++++++
.../org/apache/hop/www/async/AsyncRunServlet.java | 37 ++++++++++----------
5 files changed, 94 insertions(+), 21 deletions(-)
diff --git a/core/src/main/java/org/apache/hop/core/variables/Variables.java
b/core/src/main/java/org/apache/hop/core/variables/Variables.java
index 9dd0866fd0..b04d37b3ce 100644
--- a/core/src/main/java/org/apache/hop/core/variables/Variables.java
+++ b/core/src/main/java/org/apache/hop/core/variables/Variables.java
@@ -65,6 +65,9 @@ public class Variables implements IVariables {
// the same object as the argument.
String[] variableNames = variables.getVariableNames();
for (String variableName : variableNames) {
+ if (Utils.isEmpty(variableName)) {
+ continue;
+ }
properties.put(variableName, variables.getVariable(variableName));
}
}
@@ -140,6 +143,11 @@ public class Variables implements IVariables {
@Override
public synchronized void setVariable(String variableName, String
variableValue) {
+ // Reject null/empty names: HashMap allows null keys, but callers that
iterate
+ // names and check Immutable Sets (Set.of) throw NPE on contains(null) —
issue #7067.
+ if (Utils.isEmpty(variableName)) {
+ return;
+ }
if (variableValue != null) {
properties.put(variableName, variableValue);
} else {
diff --git
a/core/src/test/java/org/apache/hop/core/variables/VariablesTest.java
b/core/src/test/java/org/apache/hop/core/variables/VariablesTest.java
index b6e82312c5..646dfa6654 100644
--- a/core/src/test/java/org/apache/hop/core/variables/VariablesTest.java
+++ b/core/src/test/java/org/apache/hop/core/variables/VariablesTest.java
@@ -153,4 +153,23 @@ class VariablesTest {
new String[] {"DataOne", "TheDataOne"},
vars.resolve(new String[] {"${VarOne}", "The${VarOne}"}));
}
+
+ /**
+ * Null or empty variable names must not enter the properties map (issue
#7067). A null key caused
+ * NPE later when checking Const.INTERNAL_*_VARIABLES Set.of collections.
+ */
+ @Test
+ void setVariableIgnoresNullAndEmptyNames() {
+ Variables vars = new Variables();
+ vars.setVariable(null, "shouldNotStore");
+ vars.setVariable("", "shouldNotStore");
+ vars.setVariable("valid", "ok");
+
+ for (String name : vars.getVariableNames()) {
+ assertTrue(name != null && !name.isEmpty());
+ }
+ assertEquals("ok", vars.getVariable("valid"));
+ assertNull(vars.getVariable(null));
+ assertNull(vars.getVariable(""));
+ }
}
diff --git
a/engine/src/main/java/org/apache/hop/pipeline/TransformWithMappingMeta.java
b/engine/src/main/java/org/apache/hop/pipeline/TransformWithMappingMeta.java
index 96f9ebd906..ffc36cca70 100644
--- a/engine/src/main/java/org/apache/hop/pipeline/TransformWithMappingMeta.java
+++ b/engine/src/main/java/org/apache/hop/pipeline/TransformWithMappingMeta.java
@@ -313,6 +313,10 @@ public abstract class TransformWithMappingMeta<Main
extends ITransform, Data ext
}
String[] variableNames = toSpace.getVariableNames();
for (String variable : variableNames) {
+ // Skip null/empty names: Set.contains(null) throws NPE and empty keys
are invalid
+ if (Utils.isEmpty(variable)) {
+ continue;
+ }
if (fromSpace.getVariable(variable) == null) {
fromSpace.setVariable(variable, toSpace.getVariable(variable));
}
@@ -326,6 +330,10 @@ public abstract class TransformWithMappingMeta<Main
extends ITransform, Data ext
}
String[] variableNames = replaceBy.getVariableNames();
for (String variableName : variableNames) {
+ // Skip null/empty names: Immutable Set.contains(null) throws NPE (issue
#7067)
+ if (Utils.isEmpty(variableName)) {
+ continue;
+ }
if (childPipelineMeta.getVariable(variableName) != null
&& !isInternalVariable(variableName, type)) {
childPipelineMeta.setVariable(variableName,
replaceBy.getVariable(variableName));
@@ -347,10 +355,10 @@ public abstract class TransformWithMappingMeta<Main
extends ITransform, Data ext
}
private static boolean isPipelineInternalVariable(String variableName) {
- return Const.INTERNAL_PIPELINE_VARIABLES.contains(variableName);
+ return variableName != null &&
Const.INTERNAL_PIPELINE_VARIABLES.contains(variableName);
}
private static boolean isWorkflowInternalVariable(String variableName) {
- return Const.INTERNAL_WORKFLOW_VARIABLES.contains(variableName);
+ return variableName != null &&
Const.INTERNAL_WORKFLOW_VARIABLES.contains(variableName);
}
}
diff --git
a/engine/src/test/java/org/apache/hop/pipeline/TransformWithMappingMetaTest.java
b/engine/src/test/java/org/apache/hop/pipeline/TransformWithMappingMetaTest.java
index 2ac79eb245..8358f3102c 100644
---
a/engine/src/test/java/org/apache/hop/pipeline/TransformWithMappingMetaTest.java
+++
b/engine/src/test/java/org/apache/hop/pipeline/TransformWithMappingMetaTest.java
@@ -16,7 +16,11 @@
*/
package org.apache.hop.pipeline;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
import org.apache.hop.core.Const;
import org.apache.hop.core.HopEnvironment;
@@ -168,4 +172,39 @@ class TransformWithMappingMetaTest {
// keep child only variables
assertEquals(variableChildOnly,
childVariables.getVariable(variableChildOnly));
}
+
+ /**
+ * Reproduces issue #7067: a null variable name in the parent space used to
NPE when checking
+ * internal variable sets (Set.of(...).contains(null)). Async web services
could inject a null key
+ * when header content variable was unset.
+ */
+ @Test
+ void replaceVariableValuesSkipsNullVariableNames() {
+ String variableOverwrite = "paramOverwrite";
+ IVariables childVariables = new Variables();
+ childVariables.setVariable(variableOverwrite, "childValue");
+
+ IVariables replaceByParentVariables = mock(IVariables.class);
+ when(replaceByParentVariables.getVariableNames())
+ .thenReturn(new String[] {null, "", variableOverwrite});
+
when(replaceByParentVariables.getVariable(variableOverwrite)).thenReturn("parentValue");
+
+ assertDoesNotThrow(
+ () ->
+ TransformWithMappingMeta.replaceVariableValues(
+ childVariables, replaceByParentVariables));
+ assertEquals("parentValue", childVariables.getVariable(variableOverwrite));
+ }
+
+ @Test
+ void addMissingVariablesSkipsNullVariableNames() {
+ IVariables fromSpace = new Variables();
+ IVariables toSpace = mock(IVariables.class);
+ when(toSpace.getVariableNames()).thenReturn(new String[] {null, "",
"addedVar"});
+ when(toSpace.getVariable("addedVar")).thenReturn("addedValue");
+
+ assertDoesNotThrow(() ->
TransformWithMappingMeta.addMissingVariables(fromSpace, toSpace));
+ assertEquals("addedValue", fromSpace.getVariable("addedVar"));
+ assertNull(fromSpace.getVariable(null));
+ }
}
diff --git
a/plugins/misc/async/src/main/java/org/apache/hop/www/async/AsyncRunServlet.java
b/plugins/misc/async/src/main/java/org/apache/hop/www/async/AsyncRunServlet.java
index 177c5a078a..ed574ba779 100644
---
a/plugins/misc/async/src/main/java/org/apache/hop/www/async/AsyncRunServlet.java
+++
b/plugins/misc/async/src/main/java/org/apache/hop/www/async/AsyncRunServlet.java
@@ -171,11 +171,13 @@ public class AsyncRunServlet extends BaseHttpServlet
implements IHopServerPlugin
workflow.initializeFrom(variables);
workflow.setVariable("SERVER_OBJECT_ID", serverObjectId);
- // See if we need to pass a variable with the content in it...
- //
- // Read the content posted?
+ // Pass body and header content into variables when configured.
+ // Guard both independently so an unset header variable never injects a
null key
+ // (which caused NPE in Kafka Consumer init via replaceVariableValues —
issue #7067).
//
String contentVariable =
variables.resolve(webService.getBodyContentVariable());
+ String headerContentVariable =
variables.resolve(webService.getHeaderContentVariable());
+
String content = "";
if (StringUtils.isNotEmpty(contentVariable)) {
try (InputStream in = request.getInputStream()) {
@@ -193,24 +195,21 @@ public class AsyncRunServlet extends BaseHttpServlet
implements IHopServerPlugin
}
}
workflow.setVariable(contentVariable, Const.NVL(content, ""));
+ }
- String headerContentVariable =
variables.resolve(webService.getHeaderContentVariable());
- String headerContent = "";
- if (StringUtils.isNotEmpty(headerContentVariable)) {
- // Create JSON object containing all request headers
- ObjectMapper objectMapper = new ObjectMapper();
- ObjectNode headersJson = objectMapper.createObjectNode();
-
- Enumeration<String> headerNames = request.getHeaderNames();
- while (headerNames.hasMoreElements()) {
- String headerName = headerNames.nextElement();
- String headerValue = request.getHeader(headerName);
- headersJson.put(headerName, headerValue);
- }
- headerContent = objectMapper.writeValueAsString(headersJson);
- }
+ if (StringUtils.isNotEmpty(headerContentVariable)) {
+ // Create JSON object containing all request headers
+ ObjectMapper objectMapper = new ObjectMapper();
+ ObjectNode headersJson = objectMapper.createObjectNode();
- workflow.setVariable(headerContentVariable, headerContent);
+ Enumeration<String> headerNames = request.getHeaderNames();
+ while (headerNames.hasMoreElements()) {
+ String headerName = headerNames.nextElement();
+ String headerValue = request.getHeader(headerName);
+ headersJson.put(headerName, headerValue);
+ }
+ String headerContent = objectMapper.writeValueAsString(headersJson);
+ workflow.setVariable(headerContentVariable, Const.NVL(headerContent,
""));
}
// Set all the other parameters as variables/parameters...