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

oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel-kafka-connector.git


The following commit(s) were added to refs/heads/main by this push:
     new 225a5b7ab3 Fix #1795: use unmarshal instead of marshal in ckcUnMarshal 
route template (#1796)
225a5b7ab3 is described below

commit 225a5b7ab3326bd29c0122b608b248e285ced547
Author: Aditya Nikam <[email protected]>
AuthorDate: Mon Aug 24 15:44:12 2026 +0530

    Fix #1795: use unmarshal instead of marshal in ckcUnMarshal route template 
(#1796)
    
    Signed-off-by: adityaanikam <[email protected]>
---
 .../utils/CamelKafkaConnectMain.java               |  2 +-
 .../camel/kafkaconnector/DataFormatTest.java       | 30 ++++++++++++++++++++++
 2 files changed, 31 insertions(+), 1 deletion(-)

diff --git 
a/core/src/main/java/org/apache/camel/kafkaconnector/utils/CamelKafkaConnectMain.java
 
b/core/src/main/java/org/apache/camel/kafkaconnector/utils/CamelKafkaConnectMain.java
index cd7e213db8..c4782298bc 100644
--- 
a/core/src/main/java/org/apache/camel/kafkaconnector/utils/CamelKafkaConnectMain.java
+++ 
b/core/src/main/java/org/apache/camel/kafkaconnector/utils/CamelKafkaConnectMain.java
@@ -316,7 +316,7 @@ public class CamelKafkaConnectMain extends SimpleMain {
                     routeTemplate("ckcUnMarshal")
                             .templateParameter("unmarshal", "dummyDataformat")
                             .from("kamelet:source")
-                            .marshal("{{unmarshal}}")
+                            .unmarshal("{{unmarshal}}")
                             .to("kamelet:sink");
 
                     //create aggregator template
diff --git 
a/core/src/test/java/org/apache/camel/kafkaconnector/DataFormatTest.java 
b/core/src/test/java/org/apache/camel/kafkaconnector/DataFormatTest.java
index 36a886c030..240a6ec0f0 100644
--- a/core/src/test/java/org/apache/camel/kafkaconnector/DataFormatTest.java
+++ b/core/src/test/java/org/apache/camel/kafkaconnector/DataFormatTest.java
@@ -23,9 +23,12 @@ import org.apache.camel.component.hl7.HL7DataFormat;
 import org.apache.camel.component.syslog.SyslogDataFormat;
 import org.apache.camel.impl.DefaultCamelContext;
 import org.apache.camel.kafkaconnector.utils.CamelKafkaConnectMain;
+import org.apache.camel.ProducerTemplate;
+import org.apache.camel.component.syslog.SyslogMessage;
 import org.apache.kafka.connect.errors.ConnectException;
 import org.junit.jupiter.api.Test;
 
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -101,6 +104,33 @@ public class DataFormatTest {
         cms.stop();
     }
 
+    @Test
+    public void testUnmarshalRouteActuallyUnmarshals() throws Exception {
+        // https://github.com/apache/camel-kafka-connector/issues/1795
+        // ckcUnMarshal previously called .marshal() instead of .unmarshal(),
+        // so the body never actually got converted to the target type.
+        Map<String, String> props = new HashMap<>();
+        props.put("camel.source.url", "direct://test");
+        props.put("topics", "mytopic");
+        props.put("camel.source.unmarshal", "syslog");
+
+        DefaultCamelContext dcc = new DefaultCamelContext();
+        CamelKafkaConnectMain cms = 
CamelKafkaConnectMain.builder("direct://start", "log://test")
+            .withProperties(props)
+            .withUnmarshallDataFormat("syslog")
+            .build(dcc);
+
+        SyslogDataFormat syslogDf = new SyslogDataFormat();
+        dcc.getRegistry().bind("syslog", syslogDf);
+
+        cms.start();
+        ProducerTemplate template = dcc.createProducerTemplate();
+        String rawSyslogLine = "<34>Oct 11 22:14:15 mymachine su: 'su root' 
failed for lonvick on /dev/pts/8";
+        Object result = template.requestBody("direct://start", rawSyslogLine);
+        assertInstanceOf(SyslogMessage.class, result);
+        cms.stop();
+    }
+
     @Test
     public void testDataFormatLookUpInRegistry() throws Exception {
         Map<String, String> props = new HashMap<>();

Reply via email to