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 =