From dba4a407ee0e871fde9aae514c69859d3a72ad36 Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Sun, 10 Jan 2021 17:52:44 +0000 Subject: [PATCH 1/9] Fix for https://issues.apache.org/jira/browse/KAFKA-12170 Cast SMT transformation for bytes -> string. Without this fix, the conversion becomes ByteBuffer.toString(), which always gives this useless result: "java.nio.HeapByteBuffer[pos=0 lim=4 cap=4]" With this change, the byte array is converted into a hex string of the byte buffer content, for example "FEDCBA9876543210" Completed with test case and successfully tried out in a real database conversion. --- .../apache/kafka/connect/transforms/Cast.java | 13 ++++++++++++- .../kafka/connect/transforms/CastTest.java | 18 ++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index e872b336e8573..5280dbbe26754 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -39,6 +39,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.nio.ByteBuffer; import java.util.EnumSet; import java.util.HashMap; import java.util.List; @@ -364,7 +365,17 @@ private static String castToString(Object value) { if (value instanceof java.util.Date) { java.util.Date dateValue = (java.util.Date) value; return Values.dateFormatFor(dateValue).format(dateValue); - } else { + } + else if (value instanceof ByteBuffer) { + ByteBuffer byteBuffer = (ByteBuffer) value; + + StringBuilder sbuf = new StringBuilder(); + for (byte b : byteBuffer.array()) { + sbuf.append(String.format("%02X", b)); + } + return sbuf.toString(); + } + else { return value.toString(); } } diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java index 9c09554dc6875..09319ebb99c2b 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.connect.transforms; +import java.nio.ByteBuffer; import java.util.Arrays; import java.util.List; import java.util.concurrent.TimeUnit; @@ -425,6 +426,18 @@ public void castLogicalToString() { @Test public void castFieldsWithSchema() { Date day = new Date(MILLIS_PER_DAY); + + ByteBuffer byteBuffer = ByteBuffer.allocate(8); + byteBuffer.put((byte) 0xFE); + byteBuffer.put((byte) 0xDC); + byteBuffer.put((byte) 0xBA); + byteBuffer.put((byte) 0x98); + byteBuffer.put((byte) 0x76); + byteBuffer.put((byte) 0x54); + byteBuffer.put((byte) 0x32); + byteBuffer.put((byte) 0x10); + byteBuffer.flip(); + xformValue.configure(Collections.singletonMap(Cast.SPEC_CONFIG, "int8:int16,int16:int32,int32:int64,int64:boolean,float32:float64,float64:boolean,boolean:int8,string:int32,bigdecimal:string,date:string,optional:int32")); // Include an optional fields and fields with defaults to validate their values are passed through properly @@ -442,6 +455,7 @@ public void castFieldsWithSchema() { builder.field("date", org.apache.kafka.connect.data.Date.SCHEMA); builder.field("optional", Schema.OPTIONAL_FLOAT32_SCHEMA); builder.field("timestamp", Timestamp.SCHEMA); + builder.field("bytes", Schema.BYTES_SCHEMA); Schema supportedTypesSchema = builder.build(); Struct recordValue = new Struct(supportedTypesSchema); @@ -456,6 +470,8 @@ public void castFieldsWithSchema() { recordValue.put("date", day); recordValue.put("string", "42"); recordValue.put("timestamp", new Date(0)); + recordValue.put("bytes", byteBuffer); + // optional field intentionally omitted SourceRecord transformed = xformValue.apply(new SourceRecord(null, null, "topic", 0, @@ -475,6 +491,7 @@ public void castFieldsWithSchema() { assertEquals("42", ((Struct) transformed.value()).get("bigdecimal")); assertEquals(Values.dateFormatFor(day).format(day), ((Struct) transformed.value()).get("date")); assertEquals(new Date(0), ((Struct) transformed.value()).get("timestamp")); + assertEquals("FEDCBA9876543210", ((Struct) transformed.value()).get("bytes")); assertNull(((Struct) transformed.value()).get("optional")); Schema transformedSchema = ((Struct) transformed.value()).schema(); @@ -489,6 +506,7 @@ public void castFieldsWithSchema() { assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("bigdecimal").schema().type()); assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("date").schema().type()); assertEquals(Schema.OPTIONAL_INT32_SCHEMA.type(), transformedSchema.field("optional").schema().type()); + assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("bytes").schema().type()); // The following fields are not changed assertEquals(Timestamp.SCHEMA.type(), transformedSchema.field("timestamp").schema().type()); } From fa801de8630028b20f89015120de12d5df63f86a Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Thu, 21 Jan 2021 18:02:07 +0000 Subject: [PATCH 2/9] Fix cast byte array -> string for byte[] It turns out that not all connectors produce a HeapBuffer to store bytes. SAP Hana, specifically, just stores a simple byte array that was not captured by the cast transformer. Completed with test case. Fix for https://issues.apache.org/jira/browse/KAFKA-12170 --- .../apache/kafka/connect/transforms/Cast.java | 18 +++++++++++++----- .../kafka/connect/transforms/CastTest.java | 13 +++++++++++-- 2 files changed, 24 insertions(+), 7 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index 5280dbbe26754..de0d640c80a38 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -361,6 +361,7 @@ else if (value instanceof String) throw new DataException("Unexpected type in Cast transformation: " + value.getClass()); } + private static String castToString(Object value) { if (value instanceof java.util.Date) { java.util.Date dateValue = (java.util.Date) value; @@ -369,17 +370,24 @@ private static String castToString(Object value) { else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; - StringBuilder sbuf = new StringBuilder(); - for (byte b : byteBuffer.array()) { - sbuf.append(String.format("%02X", b)); - } - return sbuf.toString(); + return castByteArrayToString(byteBuffer.array()); + } + else if (value.getClass() == byte[].class) { + return castByteArrayToString((byte[]) value); } else { return value.toString(); } } + private static String castByteArrayToString(byte[] array) { + StringBuilder sbuf = new StringBuilder(); + for (byte b : array) { + sbuf.append(String.format("%02X", b)); + } + return sbuf.toString(); + } + protected abstract Schema operatingSchema(R record); protected abstract Object operatingValue(R record); diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java index 09319ebb99c2b..ed637b5fbf514 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java @@ -426,7 +426,6 @@ public void castLogicalToString() { @Test public void castFieldsWithSchema() { Date day = new Date(MILLIS_PER_DAY); - ByteBuffer byteBuffer = ByteBuffer.allocate(8); byteBuffer.put((byte) 0xFE); byteBuffer.put((byte) 0xDC); @@ -438,7 +437,10 @@ public void castFieldsWithSchema() { byteBuffer.put((byte) 0x10); byteBuffer.flip(); - xformValue.configure(Collections.singletonMap(Cast.SPEC_CONFIG, "int8:int16,int16:int32,int32:int64,int64:boolean,float32:float64,float64:boolean,boolean:int8,string:int32,bigdecimal:string,date:string,optional:int32")); + byte[] byteArray = Arrays.copyOf(byteBuffer.array(), byteBuffer.array().length); + + xformValue.configure(Collections.singletonMap(Cast.SPEC_CONFIG, + "int8:int16,int16:int32,int32:int64,int64:boolean,float32:float64,float64:boolean,boolean:int8,string:int32,bigdecimal:string,date:string,optional:int32,bytes:string,byteArray:string")); // Include an optional fields and fields with defaults to validate their values are passed through properly SchemaBuilder builder = SchemaBuilder.struct(); @@ -456,6 +458,8 @@ public void castFieldsWithSchema() { builder.field("optional", Schema.OPTIONAL_FLOAT32_SCHEMA); builder.field("timestamp", Timestamp.SCHEMA); builder.field("bytes", Schema.BYTES_SCHEMA); + builder.field("byteArray", Schema.BYTES_SCHEMA); + Schema supportedTypesSchema = builder.build(); Struct recordValue = new Struct(supportedTypesSchema); @@ -471,6 +475,7 @@ public void castFieldsWithSchema() { recordValue.put("string", "42"); recordValue.put("timestamp", new Date(0)); recordValue.put("bytes", byteBuffer); + recordValue.put("byteArray", byteArray); // optional field intentionally omitted @@ -492,6 +497,8 @@ public void castFieldsWithSchema() { assertEquals(Values.dateFormatFor(day).format(day), ((Struct) transformed.value()).get("date")); assertEquals(new Date(0), ((Struct) transformed.value()).get("timestamp")); assertEquals("FEDCBA9876543210", ((Struct) transformed.value()).get("bytes")); + assertEquals("FEDCBA9876543210", ((Struct) transformed.value()).get("byteArray")); + assertNull(((Struct) transformed.value()).get("optional")); Schema transformedSchema = ((Struct) transformed.value()).schema(); @@ -507,6 +514,8 @@ public void castFieldsWithSchema() { assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("date").schema().type()); assertEquals(Schema.OPTIONAL_INT32_SCHEMA.type(), transformedSchema.field("optional").schema().type()); assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("bytes").schema().type()); + assertEquals(Schema.STRING_SCHEMA.type(), transformedSchema.field("byteArray").schema().type()); + // The following fields are not changed assertEquals(Timestamp.SCHEMA.type(), transformedSchema.field("timestamp").schema().type()); } From 421312572f130bb694759e874e97430f7ffbf67b Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Mon, 25 Jan 2021 16:29:29 +0000 Subject: [PATCH 3/9] Fixes after failed style check --- .../java/org/apache/kafka/connect/transforms/Cast.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index de0d640c80a38..c1550288ad358 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -366,16 +366,13 @@ private static String castToString(Object value) { if (value instanceof java.util.Date) { java.util.Date dateValue = (java.util.Date) value; return Values.dateFormatFor(dateValue).format(dateValue); - } - else if (value instanceof ByteBuffer) { + } else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; return castByteArrayToString(byteBuffer.array()); - } - else if (value.getClass() == byte[].class) { + } else if (value.getClass() == byte[].class) { return castByteArrayToString((byte[]) value); - } - else { + } else { return value.toString(); } } From 4e744752488b8f23af8cb50bf7b6c0b841b43ef0 Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Thu, 18 Feb 2021 20:34:02 +0000 Subject: [PATCH 4/9] Changes made after code review by Mickael Maison (@mimaison) --- .../org/apache/kafka/connect/transforms/Cast.java | 4 +--- .../apache/kafka/connect/transforms/CastTest.java | 14 ++------------ 2 files changed, 3 insertions(+), 15 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index c1550288ad358..b8816648149a8 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -361,16 +361,14 @@ else if (value instanceof String) throw new DataException("Unexpected type in Cast transformation: " + value.getClass()); } - private static String castToString(Object value) { if (value instanceof java.util.Date) { java.util.Date dateValue = (java.util.Date) value; return Values.dateFormatFor(dateValue).format(dateValue); } else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; - return castByteArrayToString(byteBuffer.array()); - } else if (value.getClass() == byte[].class) { + } else if (value instanceof byte[]) { return castByteArrayToString((byte[]) value); } else { return value.toString(); diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java index ed637b5fbf514..4c70ec41965ed 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java @@ -426,18 +426,8 @@ public void castLogicalToString() { @Test public void castFieldsWithSchema() { Date day = new Date(MILLIS_PER_DAY); - ByteBuffer byteBuffer = ByteBuffer.allocate(8); - byteBuffer.put((byte) 0xFE); - byteBuffer.put((byte) 0xDC); - byteBuffer.put((byte) 0xBA); - byteBuffer.put((byte) 0x98); - byteBuffer.put((byte) 0x76); - byteBuffer.put((byte) 0x54); - byteBuffer.put((byte) 0x32); - byteBuffer.put((byte) 0x10); - byteBuffer.flip(); - - byte[] byteArray = Arrays.copyOf(byteBuffer.array(), byteBuffer.array().length); + byte[] byteArray = new byte[] {(byte) 0xFE, (byte) 0xDC, (byte) 0xBA, (byte) 0x98, 0x76, 0x54, 0x32, 0x10}; + ByteBuffer byteBuffer = ByteBuffer.wrap(Arrays.copyOf(byteArray, byteArray.length)); xformValue.configure(Collections.singletonMap(Cast.SPEC_CONFIG, "int8:int16,int16:int32,int32:int64,int64:boolean,float32:float64,float64:boolean,boolean:int8,string:int32,bigdecimal:string,date:string,optional:int32,bytes:string,byteArray:string")); From 302851191fbac0cd8663a731efcfb0705ade13fe Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Mon, 22 Feb 2021 10:08:42 +0000 Subject: [PATCH 5/9] Address code review feedback: protect case where ByteBuffer has not array. If the ByteBuffer does not have a backing array, we use get(byte[]) to extract the content. --- .../java/org/apache/kafka/connect/transforms/Cast.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index b8816648149a8..eea86c618752e 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -367,7 +367,14 @@ private static String castToString(Object value) { return Values.dateFormatFor(dateValue).format(dateValue); } else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; - return castByteArrayToString(byteBuffer.array()); + if (byteBuffer.hasArray()) { + return castByteArrayToString(byteBuffer.array()); + } + else { + byte[] array = new byte[byteBuffer.remaining()]; + byteBuffer.get(array); + return castByteArrayToString(array); + } } else if (value instanceof byte[]) { return castByteArrayToString((byte[]) value); } else { From 411ceb351d6c5ed0a786d41a176c91eee5603e4d Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Tue, 23 Feb 2021 09:37:20 +0000 Subject: [PATCH 6/9] Kafka has ready-made Utils class for extracting bytes from a ByteBuffer. Change made after code review. --- .../org/apache/kafka/connect/transforms/Cast.java | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index eea86c618752e..abcddf09687e9 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -22,6 +22,7 @@ import org.apache.kafka.common.cache.SynchronizedCache; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.common.utils.Utils; import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.data.ConnectSchema; import org.apache.kafka.connect.data.Date; @@ -367,16 +368,10 @@ private static String castToString(Object value) { return Values.dateFormatFor(dateValue).format(dateValue); } else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; - if (byteBuffer.hasArray()) { - return castByteArrayToString(byteBuffer.array()); - } - else { - byte[] array = new byte[byteBuffer.remaining()]; - byteBuffer.get(array); - return castByteArrayToString(array); - } + return castByteArrayToString(Utils.readBytes(byteBuffer)); } else if (value instanceof byte[]) { - return castByteArrayToString((byte[]) value); + byte[] rawBytes = (byte[]) value; + return castByteArrayToString(rawBytes); } else { return value.toString(); } From 99d26f4914123ca02cb06ce8884ab4dc504a5544 Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Thu, 25 Feb 2021 16:30:06 +0000 Subject: [PATCH 7/9] Changed conversion to a Base64-encoded string instead of Hex This matches the internal format used in other parts of the Connect code. --- .../apache/kafka/connect/transforms/Cast.java | 19 +++---------------- .../kafka/connect/transforms/CastTest.java | 4 ++-- 2 files changed, 5 insertions(+), 18 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index abcddf09687e9..422d8c3c2722b 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -41,12 +41,7 @@ import org.slf4j.LoggerFactory; import java.nio.ByteBuffer; -import java.util.EnumSet; -import java.util.HashMap; -import java.util.List; -import java.util.Locale; -import java.util.Map; -import java.util.Set; +import java.util.*; import static org.apache.kafka.connect.transforms.util.Requirements.requireMap; import static org.apache.kafka.connect.transforms.util.Requirements.requireStruct; @@ -368,23 +363,15 @@ private static String castToString(Object value) { return Values.dateFormatFor(dateValue).format(dateValue); } else if (value instanceof ByteBuffer) { ByteBuffer byteBuffer = (ByteBuffer) value; - return castByteArrayToString(Utils.readBytes(byteBuffer)); + return Base64.getEncoder().encodeToString(Utils.readBytes(byteBuffer)); } else if (value instanceof byte[]) { byte[] rawBytes = (byte[]) value; - return castByteArrayToString(rawBytes); + return Base64.getEncoder().encodeToString(rawBytes); } else { return value.toString(); } } - private static String castByteArrayToString(byte[] array) { - StringBuilder sbuf = new StringBuilder(); - for (byte b : array) { - sbuf.append(String.format("%02X", b)); - } - return sbuf.toString(); - } - protected abstract Schema operatingSchema(R record); protected abstract Object operatingValue(R record); diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java index 4c70ec41965ed..7e90a44f504ea 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/CastTest.java @@ -486,8 +486,8 @@ public void castFieldsWithSchema() { assertEquals("42", ((Struct) transformed.value()).get("bigdecimal")); assertEquals(Values.dateFormatFor(day).format(day), ((Struct) transformed.value()).get("date")); assertEquals(new Date(0), ((Struct) transformed.value()).get("timestamp")); - assertEquals("FEDCBA9876543210", ((Struct) transformed.value()).get("bytes")); - assertEquals("FEDCBA9876543210", ((Struct) transformed.value()).get("byteArray")); + assertEquals("/ty6mHZUMhA=", ((Struct) transformed.value()).get("bytes")); + assertEquals("/ty6mHZUMhA=", ((Struct) transformed.value()).get("byteArray")); assertNull(((Struct) transformed.value()).get("optional")); From e772b37620a7a1db231aeaa8d1a466e05e6edcaa Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Fri, 26 Feb 2021 09:49:21 +0000 Subject: [PATCH 8/9] Replaced wildcard with single class imports. Change requested by code review. --- .../java/org/apache/kafka/connect/transforms/Cast.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index 422d8c3c2722b..43674c1353e12 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -41,7 +41,13 @@ import org.slf4j.LoggerFactory; import java.nio.ByteBuffer; -import java.util.*; +import java.util.Base64; +import java.util.EnumSet; +import java.util.HashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Set; import static org.apache.kafka.connect.transforms.util.Requirements.requireMap; import static org.apache.kafka.connect.transforms.util.Requirements.requireStruct; From ff3239cdc3202ada25987efe2ff52224751ed1e0 Mon Sep 17 00:00:00 2001 From: Sven Erik Knop Date: Wed, 3 Mar 2021 17:26:41 +0000 Subject: [PATCH 9/9] Updated Overview and Config after code review Thank you @rhauch for your feedback. --- .../main/java/org/apache/kafka/connect/transforms/Cast.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java index 43674c1353e12..0a763cc24f85b 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Cast.java @@ -59,7 +59,8 @@ public abstract class Cast> implements Transformation // allow casting nested fields. public static final String OVERVIEW_DOC = "Cast fields or the entire key or value to a specific type, e.g. to force an integer field to a smaller " - + "width. Only simple primitive types are supported -- integers, floats, boolean, and string. " + + "width. Cast from integers, floats, boolean and string to any other type, " + + "and cast binary to string (base64 encoded)." + "

Use the concrete transformation type designed for the record key (" + Key.class.getName() + ") " + "or value (" + Value.class.getName() + ")."; @@ -85,7 +86,7 @@ public String toString() { ConfigDef.Importance.HIGH, "List of fields and the type to cast them to of the form field1:type,field2:type to cast fields of " + "Maps or Structs. A single type to cast the entire value. Valid types are int8, int16, int32, " - + "int64, float32, float64, boolean, and string."); + + "int64, float32, float64, boolean, and string. Note that binary fields can only be cast to string."); private static final String PURPOSE = "cast types";