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<>();