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 d745698614 Fix #1800: scope the SMT TypeConverter and CamelContext to 
the instance (#1811)
d745698614 is described below

commit d74569861497f374d86829f9d45d2159d4433dbb
Author: Andrea Cosentino <[email protected]>
AuthorDate: Mon Aug 24 14:56:17 2026 +0200

    Fix #1800: scope the SMT TypeConverter and CamelContext to the instance 
(#1811)
    
    CamelTypeConverterTransform declared its TypeConverter static but assigned 
it
    from the per-instance configure(), so the last SMT instance configured in a
    plugin classloader repointed the converter used by every other connector's
    transforms in that worker.
    
    Make it an instance field, and release the CamelContext that
    CamelTransformSupport creates per instance: close() was empty, so Kafka
    Connect re-instantiating transforms on every connector config update left
    those contexts accumulating for the lifetime of the worker. The core tasks
    already stop their context in stop(); this brings the SMTs in line.
    
    Also report the source and target type on a conversion failure rather than
    formatting the record value into the exception message, which Kafka Connect
    surfaces in the worker log and in the task status over the REST API.
    
    The instance-scoping test registers a converter on one instance's context 
only,
    so it fails on a static field (the second configure() repoints it) and 
passes on
    an instance field - the existing tests could not tell the two apart because 
each
    default context resolves an equivalent converter.
    
    Signed-off-by: Andrea Cosentino <[email protected]>
    Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
 .../transforms/CamelTransformSupport.java          | 17 +++++
 .../transforms/CamelTypeConverterTransform.java    |  6 +-
 .../CamelTypeConverterTransformTest.java           | 89 +++++++++++++++++++++-
 3 files changed, 109 insertions(+), 3 deletions(-)

diff --git 
a/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTransformSupport.java
 
b/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTransformSupport.java
index 4f318b0ca7..64be4ace0e 100644
--- 
a/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTransformSupport.java
+++ 
b/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTransformSupport.java
@@ -20,12 +20,29 @@ import org.apache.camel.CamelContext;
 import org.apache.camel.impl.DefaultCamelContext;
 import org.apache.kafka.connect.connector.ConnectRecord;
 import org.apache.kafka.connect.transforms.Transformation;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 public abstract class CamelTransformSupport<R extends ConnectRecord<R>> 
implements Transformation<R> {
 
+    private static final Logger LOG = 
LoggerFactory.getLogger(CamelTransformSupport.class);
+
     private final CamelContext camelContext = new DefaultCamelContext();
 
     protected CamelContext getCamelContext() {
         return camelContext;
     }
+
+    /**
+     * Stops the {@link CamelContext} created for this transform instance. 
Kafka Connect re-instantiates transforms on
+     * every connector configuration update, so a subclass must release it 
from {@link Transformation#close()} rather
+     * than let it accumulate for the lifetime of the worker.
+     */
+    protected void stopCamelContext() {
+        try {
+            camelContext.stop();
+        } catch (Exception e) {
+            LOG.warn("Failed to stop the Camel context of {}: {}", 
getClass().getSimpleName(), e.getMessage(), e);
+        }
+    }
 }
diff --git 
a/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransform.java
 
b/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransform.java
index 0ca1162c74..2078ea21fa 100644
--- 
a/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransform.java
+++ 
b/core/src/main/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransform.java
@@ -36,7 +36,7 @@ public abstract class CamelTypeConverterTransform<R extends 
ConnectRecord<R>> ex
             .define(FIELD_TARGET_TYPE_CONFIG, ConfigDef.Type.CLASS, null, 
ConfigDef.Importance.HIGH,
                     "The target field type to convert the value from, this is 
full qualified Java class, e.g: java.util.Map");
 
-    private static TypeConverter typeConverter;
+    private TypeConverter typeConverter;
     private Class<?> fieldTargetType;
 
     @Override
@@ -54,7 +54,8 @@ public abstract class CamelTypeConverterTransform<R extends 
ConnectRecord<R>> ex
         final Object convertedValue = 
typeConverter.tryConvertTo(fieldTargetType, originalValue);
 
         if (convertedValue == null) {
-            throw new DataException(String.format("CamelTypeConverter was not 
able to convert value `%s` to target type of `%s`", originalValue, 
fieldTargetType.getSimpleName()));
+            throw new DataException(String.format("CamelTypeConverter was not 
able to convert a value of type `%s` to target type of `%s`",
+                    originalValue == null ? "null" : 
originalValue.getClass().getName(), fieldTargetType.getSimpleName()));
         }
 
         return convertedValue;
@@ -80,6 +81,7 @@ public abstract class CamelTypeConverterTransform<R extends 
ConnectRecord<R>> ex
 
     @Override
     public void close() {
+        stopCamelContext();
     }
 
     @Override
diff --git 
a/core/src/test/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransformTest.java
 
b/core/src/test/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransformTest.java
index 6da72c20ab..4cc93ff110 100644
--- 
a/core/src/test/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransformTest.java
+++ 
b/core/src/test/java/org/apache/camel/kafkaconnector/transforms/CamelTypeConverterTransformTest.java
@@ -22,7 +22,10 @@ import java.util.HashMap;
 import java.util.Map;
 
 import io.netty.buffer.Unpooled;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
 import org.apache.camel.component.netty.http.NettyChannelBufferStreamCache;
+import org.apache.camel.support.TypeConverterSupport;
 import org.apache.kafka.common.config.ConfigException;
 import org.apache.kafka.connect.data.Schema;
 import org.apache.kafka.connect.data.SchemaBuilder;
@@ -31,8 +34,12 @@ import org.apache.kafka.connect.source.SourceRecord;
 import org.apache.kafka.connect.transforms.Transformation;
 import org.junit.jupiter.api.Test;
 
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -134,4 +141,84 @@ public class CamelTypeConverterTransformTest {
         assertThrows(ConfigException.class, () -> 
transformationKey.configure(Collections.emptyMap()));
     }
 
-}
\ No newline at end of file
+    @Test
+    public void testEachInstanceKeepsItsOwnCamelContext() {
+        final CamelTypeConverterTransform.Value<SourceRecord> first = new 
CamelTypeConverterTransform.Value<>();
+        final CamelTypeConverterTransform.Value<SourceRecord> second = new 
CamelTypeConverterTransform.Value<>();
+
+        assertNotSame(first.getCamelContext(), second.getCamelContext());
+    }
+
+    @Test
+    public void testConverterIsScopedToTheInstanceThatConfiguredIt() {
+        final CamelTypeConverterTransform.Value<SourceRecord> first = new 
CamelTypeConverterTransform.Value<>();
+
+        // teach ONLY this instance's context how to produce a Marker
+        
first.getCamelContext().getTypeConverterRegistry().addTypeConverter(Marker.class,
 String.class,
+                new TypeConverterSupport() {
+                    @Override
+                    public <T> T convertTo(Class<T> type, Exchange exchange, 
Object value) {
+                        return type.cast(new Marker(String.valueOf(value)));
+                    }
+                });
+
+        final Map<String, Object> toMarker = new HashMap<>();
+        toMarker.put(CamelTypeConverterTransform.FIELD_TARGET_TYPE_CONFIG, 
Marker.class.getName());
+        first.configure(toMarker);
+
+        // configuring a second instance afterwards must not repoint the 
converter the first one uses
+        final Map<String, Object> toString = new HashMap<>();
+        toString.put(CamelTypeConverterTransform.FIELD_TARGET_TYPE_CONFIG, 
String.class.getName());
+        new 
CamelTypeConverterTransform.Value<SourceRecord>().configure(toString);
+
+        final SourceRecord record = new SourceRecord(Collections.emptyMap(), 
Collections.emptyMap(), "topic",
+                Schema.STRING_SCHEMA, "1234", Schema.STRING_SCHEMA, "abc");
+
+        assertInstanceOf(Marker.class, first.apply(record).value());
+    }
+
+    /** Target type known only to the converter registered on one instance's 
context. */
+    public static final class Marker {
+        private final String value;
+
+        Marker(String value) {
+            this.value = value;
+        }
+
+        @Override
+        public String toString() {
+            return "Marker[" + value + "]";
+        }
+    }
+
+    @Test
+    public void testCloseStopsTheCamelContext() {
+        final CamelTypeConverterTransform.Value<SourceRecord> transform = new 
CamelTypeConverterTransform.Value<>();
+        final Map<String, Object> props = new HashMap<>();
+        props.put(CamelTypeConverterTransform.FIELD_TARGET_TYPE_CONFIG, 
String.class.getName());
+        transform.configure(props);
+
+        final CamelContext context = transform.getCamelContext();
+        assertDoesNotThrow(transform::close);
+        assertFalse(context.getStatus().isStarted());
+    }
+
+    @Test
+    public void testConversionFailureDoesNotEchoTheRecordValue() {
+        final Map<String, Object> props = new HashMap<>();
+        props.put(CamelTypeConverterTransform.FIELD_TARGET_TYPE_CONFIG, 
java.time.LocalDate.class.getName());
+
+        final Transformation<SourceRecord> transform = new 
CamelTypeConverterTransform.Value<>();
+        transform.configure(props);
+
+        final String secret = "s3cr3t-record-content";
+        final SourceRecord record = new SourceRecord(Collections.emptyMap(), 
Collections.emptyMap(), "topic",
+                Schema.STRING_SCHEMA, "1234", Schema.STRING_SCHEMA, secret);
+
+        final org.apache.kafka.connect.errors.DataException e =
+                
assertThrows(org.apache.kafka.connect.errors.DataException.class, () -> 
transform.apply(record));
+
+        assertFalse(e.getMessage().contains(secret), "the record value must 
not be echoed into the exception message");
+        assertTrue(e.getMessage().contains(String.class.getName()), "the 
source type should be reported instead");
+    }
+}

Reply via email to