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

bamaer 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 1c72f98927 Issue #4106 : Map imported KafkaConsumerInput steps to 
KafkaConsumer (#8686)
1c72f98927 is described below

commit 1c72f989271b6285761a27479f25a5f5bc5db20b
Author: Matt Casters <[email protected]>
AuthorDate: Thu Oct 1 11:20:09 2026 +0200

    Issue #4106 : Map imported KafkaConsumerInput steps to KafkaConsumer (#8686)
---
 .../org/apache/hop/imports/kettle/KettleConst.java |  2 +
 .../kettle/KettleImportKafkaConsumerTest.java      | 86 +++++++++++++++++++++-
 2 files changed, 87 insertions(+), 1 deletion(-)

diff --git 
a/plugins/misc/import/src/main/java/org/apache/hop/imports/kettle/KettleConst.java
 
b/plugins/misc/import/src/main/java/org/apache/hop/imports/kettle/KettleConst.java
index 234ca9357e..beb236a502 100644
--- 
a/plugins/misc/import/src/main/java/org/apache/hop/imports/kettle/KettleConst.java
+++ 
b/plugins/misc/import/src/main/java/org/apache/hop/imports/kettle/KettleConst.java
@@ -180,6 +180,8 @@ public class KettleConst {
                 // Text File Input deprecated
                 {"TextFileInput", "TextFileInput2"},
                 {"HTTP", "Http"},
+                // Pentaho's Kafka Consumer step id. The older key is kept for 
files that used it.
+                {"KafkaConsumerInput", "KafkaConsumer"},
                 {"KettleKafkaConsumerInput", "KafkaConsumer"},
                 {"PentahoGoogleSheetsPluginOutputMeta", "GoogleSheetsOutput"},
                 {"PentahoGoogleSheetsPluginInputMeta", "GoogleSheetsInput"}
diff --git 
a/plugins/misc/import/src/test/java/org/apache/hop/imports/kettle/KettleImportKafkaConsumerTest.java
 
b/plugins/misc/import/src/test/java/org/apache/hop/imports/kettle/KettleImportKafkaConsumerTest.java
index eccad96af9..d019b3c7ec 100644
--- 
a/plugins/misc/import/src/test/java/org/apache/hop/imports/kettle/KettleImportKafkaConsumerTest.java
+++ 
b/plugins/misc/import/src/test/java/org/apache/hop/imports/kettle/KettleImportKafkaConsumerTest.java
@@ -23,11 +23,14 @@ import static org.junit.jupiter.api.Assertions.assertNull;
 import java.io.ByteArrayInputStream;
 import java.lang.reflect.Method;
 import java.nio.charset.StandardCharsets;
+import org.apache.hop.core.HopClientEnvironment;
 import org.apache.hop.core.xml.XmlHandler;
 import org.apache.hop.core.xml.XmlParserFactoryProducer;
+import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.ValueSource;
 import org.w3c.dom.Document;
 import org.w3c.dom.Node;
 
@@ -35,13 +38,22 @@ import org.w3c.dom.Node;
  * A Kettle Kafka consumer stores its sub-transformation in {@code 
transformationPath}. Repository
  * references omit the {@code .ktr} extension, and Hop cannot open the 
imported pipeline until that
  * path ends with {@code .hpl} (issue #4107).
+ *
+ * <p>Pentaho writes the Kafka Consumer step as {@code KafkaConsumerInput}. 
Hop's transform id is
+ * {@code KafkaConsumer}. Leaving the Kettle id in place makes the imported 
transform unopenable
+ * (issue #4106).
  */
 class KettleImportKafkaConsumerTest {
 
   private static final String ENTRY_TYPE = 
"org.apache.hop.imports.kettle.KettleImport$EntryType";
 
+  @BeforeAll
+  static void setUpBeforeClass() throws Exception {
+    HopClientEnvironment.init();
+  }
+
   @ParameterizedTest
-  @CsvSource({"KafkaConsumerInput,KafkaConsumerInput", 
"KettleKafkaConsumerInput,KafkaConsumer"})
+  @CsvSource({"KafkaConsumerInput,KafkaConsumer", 
"KettleKafkaConsumerInput,KafkaConsumer"})
   void testPathWithoutExtensionGainsHpl(String kettleType, String hopType) 
throws Exception {
     Node transform =
         importKafkaStep(
@@ -130,6 +142,78 @@ class KettleImportKafkaConsumerTest {
         XmlHandler.getTagValue(transform, "pipelinePath"));
   }
 
+  @ParameterizedTest
+  @ValueSource(strings = {"KafkaConsumerInput", "KettleKafkaConsumerInput"})
+  void testKafkaConsumerTypeBecomesKafkaConsumer(String kettleType) throws 
Exception {
+    Document doc = parse(kettleKafkaConsumer(kettleType));
+    processNode(doc);
+
+    Node pipeline = XmlHandler.getSubNode(doc, "pipeline");
+    assertNotNull(pipeline);
+    Node transform = XmlHandler.getSubNode(pipeline, "transform");
+    assertNotNull(transform);
+    assertEquals("kafkaCnsmr:zsp1", XmlHandler.getTagValue(transform, "name"));
+    assertEquals("KafkaConsumer", XmlHandler.getTagValue(transform, "type"));
+    assertNull(XmlHandler.getSubNode(transform, "transformationPath"));
+    assertEquals(
+        "${Internal.Entry.Current.Folder}/child.hpl",
+        XmlHandler.getTagValue(transform, "pipelinePath"));
+    assertNull(XmlHandler.getSubNode(transform, "SUB_STEP"));
+    assertEquals("Output", XmlHandler.getTagValue(transform, "subTransform"));
+    assertEquals("orders", XmlHandler.getTagValue(transform, "topic"));
+    assertEquals("group-a", XmlHandler.getTagValue(transform, 
"consumerGroup"));
+    assertEquals("100", XmlHandler.getTagValue(transform, "batchSize"));
+    assertEquals("1000", XmlHandler.getTagValue(transform, "batchDuration"));
+    assertEquals("localhost:9092", XmlHandler.getTagValue(transform, 
"directBootstrapServers"));
+    assertEquals("Y", XmlHandler.getTagValue(transform, "AUTO_COMMIT"));
+
+    Node keyField = XmlHandler.getSubNodeByNr(transform, "OutputField", 0);
+    assertEquals("key", XmlHandler.getTagAttribute(keyField, "kafkaName"));
+    assertEquals("String", XmlHandler.getTagAttribute(keyField, "type"));
+    assertEquals("Key", XmlHandler.getNodeValue(keyField));
+
+    Node option =
+        XmlHandler.getSubNode(XmlHandler.getSubNode(transform, 
"advancedConfig"), "option");
+    assertEquals("auto.offset.reset", XmlHandler.getTagAttribute(option, 
"property"));
+    assertEquals("earliest", XmlHandler.getTagAttribute(option, "value"));
+  }
+
+  @Test
+  void testOtherStepTypesAreLeftAlone() throws Exception {
+    Document doc =
+        parse(
+            
"<transformation><step><name>read</name><type>TableInput</type></step></transformation>");
+    processNode(doc);
+
+    Node transform = XmlHandler.getSubNode(XmlHandler.getSubNode(doc, 
"pipeline"), "transform");
+    assertEquals("TableInput", XmlHandler.getTagValue(transform, "type"));
+  }
+
+  private static String kettleKafkaConsumer(String type) {
+    return "<transformation>"
+        + "<step>"
+        + "<name>kafkaCnsmr:zsp1</name>"
+        + "<type>"
+        + type
+        + "</type>"
+        + "<topic>orders</topic>"
+        + "<consumerGroup>group-a</consumerGroup>"
+        + 
"<transformationPath>${Internal.Entry.Current.Directory}/child.ktr</transformationPath>"
+        + "<SUB_STEP>Output</SUB_STEP>"
+        + "<batchSize>100</batchSize>"
+        + "<batchDuration>1000</batchDuration>"
+        + "<connectionType>DIRECT</connectionType>"
+        + "<directBootstrapServers>localhost:9092</directBootstrapServers>"
+        + "<AUTO_COMMIT>Y</AUTO_COMMIT>"
+        + "<OutputField kafkaName=\"key\" type=\"String\">Key</OutputField>"
+        + "<OutputField kafkaName=\"message\" 
type=\"String\">Message</OutputField>"
+        + "<advancedConfig>"
+        + "<option property=\"auto.offset.reset\" value=\"earliest\"/>"
+        + "</advancedConfig>"
+        + "</step>"
+        + "</transformation>";
+  }
+
   private Node importKafkaStep(String type, String transformationPath, boolean 
withRemoteSteps)
       throws Exception {
     String remoteSteps =

Reply via email to