From d6a5243774ef04efb927973e56001f548054ed54 Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Mon, 21 Oct 2019 13:07:15 -0500 Subject: [PATCH 1/5] =?UTF-8?q?KAFKA-9074:=20Correct=20Connect=E2=80=99s?= =?UTF-8?q?=20`Values.parseString`=20to=20properly=20parse=20a=20time=20an?= =?UTF-8?q?d=20timestamp=20literal?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Time and timestamp literal strings contain a `:` character, but the internal parser used in the `Values.parseString(String)` method tokenizes on the colon character to tokenize and parse map entries. The colon could be escaped, but then the backslash character used to escape the colon is not removed and the parser fails to match the literal as a time or timestamp value. This fix corrects the parsing logic to properly parse timestamp and time literal strings whose colon characters are either escaped or unescaped. Additional unit tests were added to first verify the incorrect behavior and then to validate the correction. --- .../org/apache/kafka/connect/data/Values.java | 99 +++++++++++++------ .../apache/kafka/connect/data/ValuesTest.java | 68 +++++++++++++ 2 files changed, 137 insertions(+), 30 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java index 93c320a236760..b56d2d8d458e3 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java @@ -70,9 +70,9 @@ public class Values { private static final String FALSE_LITERAL = Boolean.FALSE.toString(); private static final long MILLIS_PER_DAY = 24 * 60 * 60 * 1000; private static final String NULL_VALUE = "null"; - private static final String ISO_8601_DATE_FORMAT_PATTERN = "yyyy-MM-dd"; - private static final String ISO_8601_TIME_FORMAT_PATTERN = "HH:mm:ss.SSS'Z'"; - private static final String ISO_8601_TIMESTAMP_FORMAT_PATTERN = ISO_8601_DATE_FORMAT_PATTERN + "'T'" + ISO_8601_TIME_FORMAT_PATTERN; + protected static final String ISO_8601_DATE_FORMAT_PATTERN = "yyyy-MM-dd"; + protected static final String ISO_8601_TIME_FORMAT_PATTERN = "HH:mm:ss.SSS'Z'"; + protected static final String ISO_8601_TIMESTAMP_FORMAT_PATTERN = ISO_8601_DATE_FORMAT_PATTERN + "'T'" + ISO_8601_TIME_FORMAT_PATTERN; private static final String QUOTE_DELIMITER = "\""; private static final String COMMA_DELIMITER = ","; @@ -838,7 +838,9 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No } if (!parser.canConsume(ENTRY_DELIMITER)) { - throw new DataException("Map entry is missing '=': " + parser.original()); + throw new DataException("Map entry is missing '" + ENTRY_DELIMITER + + "' at " + parser.position() + + " in " + parser.original()); } SchemaAndValue value = parse(parser, true); Object entryValue = value != null ? value.value() : null; @@ -867,6 +869,24 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No char firstChar = token.charAt(0); boolean firstCharIsDigit = Character.isDigit(firstChar); + + // Temporal types are more restrictive, so try them first + if (firstCharIsDigit) { + // The time and timestamp literals may be split into 5 tokens since an unescaped colon + // is a delimiter. Check these first since the first of these tokens is a simple numeric + int position = parser.mark(); + String timeOrTimestampStr = parser.nextTokens(4, startPosition); // may be null + SchemaAndValue temporal = parseAsTemporal(timeOrTimestampStr); + if (temporal != null) { + return temporal; + } + // No match was found using the 5 tokens, so rewind and see if the current token has a date, time, or timestamp + parser.rewindTo(position); + temporal = parseAsTemporal(token); + if (temporal != null) { + return temporal; + } + } if (firstCharIsDigit || firstChar == '+' || firstChar == '-') { try { // Try to parse as a number ... @@ -901,37 +921,40 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No // can't parse as a number } } - if (firstCharIsDigit) { - // Check for a date, time, or timestamp ... - int tokenLength = token.length(); - if (tokenLength == ISO_8601_DATE_LENGTH) { - try { - return new SchemaAndValue(Date.SCHEMA, new SimpleDateFormat(ISO_8601_DATE_FORMAT_PATTERN).parse(token)); - } catch (ParseException e) { - // not a valid date - } - } else if (tokenLength == ISO_8601_TIME_LENGTH) { - try { - return new SchemaAndValue(Time.SCHEMA, new SimpleDateFormat(ISO_8601_TIME_FORMAT_PATTERN).parse(token)); - } catch (ParseException e) { - // not a valid date - } - } else if (tokenLength == ISO_8601_TIMESTAMP_LENGTH) { - try { - return new SchemaAndValue(Timestamp.SCHEMA, new SimpleDateFormat(ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(token)); - } catch (ParseException e) { - // not a valid date - } - } - } - if (embedded) { - throw new DataException("Failed to parse embedded value"); - } // At this point, the only thing this can be is a string. Embedded strings were processed above, // so this is not embedded and we can use the original string... return new SchemaAndValue(Schema.STRING_SCHEMA, parser.original()); } + protected static SchemaAndValue parseAsTemporal(String token) { + if (token == null) { + return null; + } + // If the colons were escaped, we'll see the escape chars and need to remove them + token = token.replace("\\:", ":"); + int tokenLength = token.length(); + if (tokenLength == ISO_8601_TIME_LENGTH) { + try { + return new SchemaAndValue(Time.SCHEMA, new SimpleDateFormat(ISO_8601_TIME_FORMAT_PATTERN).parse(token)); + } catch (ParseException e) { + // not a valid date + } + } else if (tokenLength == ISO_8601_TIMESTAMP_LENGTH) { + try { + return new SchemaAndValue(Timestamp.SCHEMA, new SimpleDateFormat(ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(token)); + } catch (ParseException e) { + // not a valid date + } + } else if (tokenLength == ISO_8601_DATE_LENGTH) { + try { + return new SchemaAndValue(Date.SCHEMA, new SimpleDateFormat(ISO_8601_DATE_FORMAT_PATTERN).parse(token)); + } catch (ParseException e) { + // not a valid date + } + } + return null; + } + protected static Schema commonSchemaFor(Schema previous, SchemaAndValue latest) { if (latest == null) { return previous; @@ -1112,6 +1135,22 @@ public String next() { return previousToken; } + public String nextTokens(int n) { + return nextTokens(n, mark()); + } + + public String nextTokens(int n, int startingIndex) { + int start = mark(); + for (int i=0; i!=n; ++i) { + if (!hasNext()) { + rewindTo(start); + return null; + } + next(); + } + return original.substring(startingIndex, position()); + } + private String consumeNextToken() throws NoSuchElementException { boolean escaped = false; int start = iter.getIndex(); diff --git a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java index a5909f3ece9bf..cf6edc2a724c3 100644 --- a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java +++ b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java @@ -21,6 +21,8 @@ import org.apache.kafka.connect.errors.DataException; import org.junit.Test; +import java.text.SimpleDateFormat; +import java.time.Instant; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; @@ -383,6 +385,58 @@ public void shouldParseStringListWithExtraDelimitersAndReturnString() { assertEquals(str, result.value()); } + @Test + public void shouldParseTimestampStringAsTimestamp() throws Exception { + String str = "2019-08-23T14:34:54.346Z"; + SchemaAndValue result = Values.parseString(str); + assertEquals(Type.INT64, result.schema().type()); + assertEquals(Timestamp.LOGICAL_NAME, result.schema().name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(str); + assertEquals(expected, result.value()); + } + + @Test + public void shouldParseDateStringAsDate() throws Exception { + String str = "2019-08-23"; + SchemaAndValue result = Values.parseString(str); + assertEquals(Type.INT32, result.schema().type()); + assertEquals(Date.LOGICAL_NAME, result.schema().name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_DATE_FORMAT_PATTERN).parse(str); + assertEquals(expected, result.value()); + } + + @Test + public void shouldParseTimeStringAsDate() throws Exception { + String str = "14:34:54.346Z"; + SchemaAndValue result = Values.parseString(str); + assertEquals(Type.INT32, result.schema().type()); + assertEquals(Time.LOGICAL_NAME, result.schema().name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIME_FORMAT_PATTERN).parse(str); + assertEquals(expected, result.value()); + } + + @Test + public void shouldParseTimestampStringWithEscapedColonsAsTimestamp() throws Exception { + String str = "2019-08-23T14\\:34\\:54.346Z"; + SchemaAndValue result = Values.parseString(str); + assertEquals(Type.INT64, result.schema().type()); + assertEquals(Timestamp.LOGICAL_NAME, result.schema().name()); + String expectedStr = "2019-08-23T14:34:54.346Z"; + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(expectedStr); + assertEquals(expected, result.value()); + } + + @Test + public void shouldParseTimeStringWithEscapedColonsAsDate() throws Exception { + String str = "14\\:34\\:54.346Z"; + SchemaAndValue result = Values.parseString(str); + assertEquals(Type.INT32, result.schema().type()); + assertEquals(Time.LOGICAL_NAME, result.schema().name()); + String expectedStr = "14:34:54.346Z"; + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIME_FORMAT_PATTERN).parse(expectedStr); + assertEquals(expected, result.value()); + } + /** * This is technically invalid JSON, and we don't want to simply ignore the blank elements. */ @@ -432,6 +486,20 @@ public void shouldFailToParseStringOfMapWithIntValuesWithBlankEntries() { Values.convertToMap(Schema.STRING_SCHEMA, " { \"foo\" : \"1234567890\" ,, \"bar\" : \"0\", \"baz\" : \"boz\" } "); } + @Test + public void shouldConsumeMultipleTokens() { + String value = "a:b:c:d:e:f:g:h"; + Parser parser = new Parser(value); + String firstFive = parser.nextTokens(5); + assertEquals("a:b:c", firstFive); + assertEquals(":", parser.next()); + assertEquals("d", parser.next()); + assertEquals(":", parser.next()); + String lastEight = parser.nextTokens(8); // only 7 remain + assertNull(lastEight); + assertEquals("e", parser.next()); + } + @Test public void shouldParseStringsWithoutDelimiters() { //assertParsed(""); From 39922d89b66bd20a45302fc108ed61c5ffe5b51f Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Mon, 21 Oct 2019 13:28:28 -0500 Subject: [PATCH 2/5] Fixed checkstyle problems --- .../src/main/java/org/apache/kafka/connect/data/Values.java | 4 ++-- .../test/java/org/apache/kafka/connect/data/ValuesTest.java | 1 - 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java index b56d2d8d458e3..0fe5815b34612 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java @@ -941,7 +941,7 @@ protected static SchemaAndValue parseAsTemporal(String token) { } } else if (tokenLength == ISO_8601_TIMESTAMP_LENGTH) { try { - return new SchemaAndValue(Timestamp.SCHEMA, new SimpleDateFormat(ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(token)); + return new SchemaAndValue(Timestamp.SCHEMA, new SimpleDateFormat(ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(token)); } catch (ParseException e) { // not a valid date } @@ -1141,7 +1141,7 @@ public String nextTokens(int n) { public String nextTokens(int n, int startingIndex) { int start = mark(); - for (int i=0; i!=n; ++i) { + for (int i = 0; i != n; ++i) { if (!hasNext()) { rewindTo(start); return null; diff --git a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java index cf6edc2a724c3..c95b5a46c7c86 100644 --- a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java +++ b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java @@ -22,7 +22,6 @@ import org.junit.Test; import java.text.SimpleDateFormat; -import java.time.Instant; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; From e22892d1ade5160b3c5113ce515255763eb6749f Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Tue, 22 Oct 2019 10:08:43 -0500 Subject: [PATCH 3/5] Incorporated feedback --- .../org/apache/kafka/connect/data/Values.java | 30 +++++++++---------- .../apache/kafka/connect/data/ValuesTest.java | 4 +-- 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java index 0fe5815b34612..3c263845edb54 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java @@ -70,9 +70,9 @@ public class Values { private static final String FALSE_LITERAL = Boolean.FALSE.toString(); private static final long MILLIS_PER_DAY = 24 * 60 * 60 * 1000; private static final String NULL_VALUE = "null"; - protected static final String ISO_8601_DATE_FORMAT_PATTERN = "yyyy-MM-dd"; - protected static final String ISO_8601_TIME_FORMAT_PATTERN = "HH:mm:ss.SSS'Z'"; - protected static final String ISO_8601_TIMESTAMP_FORMAT_PATTERN = ISO_8601_DATE_FORMAT_PATTERN + "'T'" + ISO_8601_TIME_FORMAT_PATTERN; + static final String ISO_8601_DATE_FORMAT_PATTERN = "yyyy-MM-dd"; + static final String ISO_8601_TIME_FORMAT_PATTERN = "HH:mm:ss.SSS'Z'"; + static final String ISO_8601_TIMESTAMP_FORMAT_PATTERN = ISO_8601_DATE_FORMAT_PATTERN + "'T'" + ISO_8601_TIME_FORMAT_PATTERN; private static final String QUOTE_DELIMITER = "\""; private static final String COMMA_DELIMITER = ","; @@ -875,14 +875,17 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No // The time and timestamp literals may be split into 5 tokens since an unescaped colon // is a delimiter. Check these first since the first of these tokens is a simple numeric int position = parser.mark(); - String timeOrTimestampStr = parser.nextTokens(4, startPosition); // may be null - SchemaAndValue temporal = parseAsTemporal(timeOrTimestampStr); - if (temporal != null) { - return temporal; + String remainder = parser.next(4); + if (remainder != null) { + String timeOrTimestampStr = token + remainder; + SchemaAndValue temporal = parseAsTemporal(timeOrTimestampStr); + if (temporal != null) { + return temporal; + } } // No match was found using the 5 tokens, so rewind and see if the current token has a date, time, or timestamp parser.rewindTo(position); - temporal = parseAsTemporal(token); + SchemaAndValue temporal = parseAsTemporal(token); if (temporal != null) { return temporal; } @@ -926,7 +929,7 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No return new SchemaAndValue(Schema.STRING_SCHEMA, parser.original()); } - protected static SchemaAndValue parseAsTemporal(String token) { + private static SchemaAndValue parseAsTemporal(String token) { if (token == null) { return null; } @@ -1135,11 +1138,8 @@ public String next() { return previousToken; } - public String nextTokens(int n) { - return nextTokens(n, mark()); - } - - public String nextTokens(int n, int startingIndex) { + public String next(int n) { + int current = mark(); int start = mark(); for (int i = 0; i != n; ++i) { if (!hasNext()) { @@ -1148,7 +1148,7 @@ public String nextTokens(int n, int startingIndex) { } next(); } - return original.substring(startingIndex, position()); + return original.substring(current, position()); } private String consumeNextToken() throws NoSuchElementException { diff --git a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java index c95b5a46c7c86..85c822cf13491 100644 --- a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java +++ b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java @@ -489,12 +489,12 @@ public void shouldFailToParseStringOfMapWithIntValuesWithBlankEntries() { public void shouldConsumeMultipleTokens() { String value = "a:b:c:d:e:f:g:h"; Parser parser = new Parser(value); - String firstFive = parser.nextTokens(5); + String firstFive = parser.next(5); assertEquals("a:b:c", firstFive); assertEquals(":", parser.next()); assertEquals("d", parser.next()); assertEquals(":", parser.next()); - String lastEight = parser.nextTokens(8); // only 7 remain + String lastEight = parser.next(8); // only 7 remain assertNull(lastEight); assertEquals("e", parser.next()); } From 435b15f71bec4b09c04a607703f00e8d3c6e714c Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Wed, 23 Oct 2019 09:12:52 -0500 Subject: [PATCH 4/5] KAFKA-9074: Corrected the parsing of logical types within arrays and maps --- .../org/apache/kafka/connect/data/Values.java | 32 ++++++- .../apache/kafka/connect/data/ValuesTest.java | 88 +++++++++++++++++++ 2 files changed, 119 insertions(+), 1 deletion(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java index 3c263845edb54..bd6fc217d541f 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java @@ -30,13 +30,17 @@ import java.text.SimpleDateFormat; import java.text.StringCharacterIterator; import java.util.ArrayList; +import java.util.Arrays; import java.util.Base64; import java.util.Calendar; +import java.util.Collections; +import java.util.HashSet; import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.NoSuchElementException; +import java.util.Set; import java.util.TimeZone; import java.util.regex.Pattern; @@ -73,6 +77,15 @@ public class Values { static final String ISO_8601_DATE_FORMAT_PATTERN = "yyyy-MM-dd"; static final String ISO_8601_TIME_FORMAT_PATTERN = "HH:mm:ss.SSS'Z'"; static final String ISO_8601_TIMESTAMP_FORMAT_PATTERN = ISO_8601_DATE_FORMAT_PATTERN + "'T'" + ISO_8601_TIME_FORMAT_PATTERN; + private static final Set TEMPORAL_LOGICAL_TYPE_NAMES = + Collections.unmodifiableSet( + new HashSet<>( + Arrays.asList(Time.LOGICAL_NAME, + Timestamp.LOGICAL_NAME, + Date.LOGICAL_NAME + ) + ) + ); private static final String QUOTE_DELIMITER = "\""; private static final String COMMA_DELIMITER = ","; @@ -467,6 +480,9 @@ protected static Object convertTo(Schema toSchema, Schema fromSchema, Object val int days = (int) (millis / MILLIS_PER_DAY); // truncates return Date.toLogical(toSchema, days); } + } else { + // There is no fromSchema, so no conversion is needed + return value; } } long numeric = asLong(value, fromSchema, null); @@ -492,6 +508,9 @@ protected static Object convertTo(Schema toSchema, Schema fromSchema, Object val calendar.set(Calendar.DAY_OF_MONTH, 1); return Time.toLogical(toSchema, (int) calendar.getTimeInMillis()); } + } else { + // There is no fromSchema, so no conversion is needed + return value; } } long numeric = asLong(value, fromSchema, null); @@ -523,6 +542,9 @@ protected static Object convertTo(Schema toSchema, Schema fromSchema, Object val if (Timestamp.LOGICAL_NAME.equals(fromSchemaName)) { return value; } + } else { + // There is no fromSchema, so no conversion is needed + return value; } } long numeric = asLong(value, fromSchema, null); @@ -755,7 +777,14 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No } sb.append(parser.next()); } - return new SchemaAndValue(Schema.STRING_SCHEMA, sb.toString()); + String content = sb.toString(); + // We can parse string literals as temporal logical types, but all others + // are treated as strings + SchemaAndValue parsed = parseString(content); + if (parsed != null && TEMPORAL_LOGICAL_TYPE_NAMES.contains(parsed.schema().name())) { + return parsed; + } + return new SchemaAndValue(Schema.STRING_SCHEMA, content); } } @@ -1114,6 +1143,7 @@ public int mark() { public void rewindTo(int position) { iter.setIndex(position); nextToken = null; + previousToken = null; } public String original() { diff --git a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java index 85c822cf13491..c437e46c25956 100644 --- a/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java +++ b/connect/api/src/test/java/org/apache/kafka/connect/data/ValuesTest.java @@ -436,6 +436,94 @@ public void shouldParseTimeStringWithEscapedColonsAsDate() throws Exception { assertEquals(expected, result.value()); } + @Test + public void shouldParseDateStringAsDateInArray() throws Exception { + String dateStr = "2019-08-23"; + String arrayStr = "[" + dateStr + "]"; + SchemaAndValue result = Values.parseString(arrayStr); + assertEquals(Type.ARRAY, result.schema().type()); + Schema elementSchema = result.schema().valueSchema(); + assertEquals(Type.INT32, elementSchema.type()); + assertEquals(Date.LOGICAL_NAME, elementSchema.name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_DATE_FORMAT_PATTERN).parse(dateStr); + assertEquals(Collections.singletonList(expected), result.value()); + } + + @Test + public void shouldParseTimeStringAsTimeInArray() throws Exception { + String timeStr = "14:34:54.346Z"; + String arrayStr = "[" + timeStr + "]"; + SchemaAndValue result = Values.parseString(arrayStr); + assertEquals(Type.ARRAY, result.schema().type()); + Schema elementSchema = result.schema().valueSchema(); + assertEquals(Type.INT32, elementSchema.type()); + assertEquals(Time.LOGICAL_NAME, elementSchema.name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIME_FORMAT_PATTERN).parse(timeStr); + assertEquals(Collections.singletonList(expected), result.value()); + } + + @Test + public void shouldParseTimestampStringAsTimestampInArray() throws Exception { + String tsStr = "2019-08-23T14:34:54.346Z"; + String arrayStr = "[" + tsStr + "]"; + SchemaAndValue result = Values.parseString(arrayStr); + assertEquals(Type.ARRAY, result.schema().type()); + Schema elementSchema = result.schema().valueSchema(); + assertEquals(Type.INT64, elementSchema.type()); + assertEquals(Timestamp.LOGICAL_NAME, elementSchema.name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(tsStr); + assertEquals(Collections.singletonList(expected), result.value()); + } + + @Test + public void shouldParseMultipleTimestampStringAsTimestampInArray() throws Exception { + String tsStr1 = "2019-08-23T14:34:54.346Z"; + String tsStr2 = "2019-01-23T15:12:34.567Z"; + String tsStr3 = "2019-04-23T19:12:34.567Z"; + String arrayStr = "[" + tsStr1 + "," + tsStr2 + ", " + tsStr3 + "]"; + SchemaAndValue result = Values.parseString(arrayStr); + assertEquals(Type.ARRAY, result.schema().type()); + Schema elementSchema = result.schema().valueSchema(); + assertEquals(Type.INT64, elementSchema.type()); + assertEquals(Timestamp.LOGICAL_NAME, elementSchema.name()); + java.util.Date expected1 = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(tsStr1); + java.util.Date expected2 = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(tsStr2); + java.util.Date expected3 = new SimpleDateFormat(Values.ISO_8601_TIMESTAMP_FORMAT_PATTERN).parse(tsStr3); + assertEquals(Arrays.asList(expected1, expected2, expected3), result.value()); + } + + @Test + public void shouldParseQuotedTimeStringAsTimeInMap() throws Exception { + String keyStr = "k1"; + String timeStr = "14:34:54.346Z"; + String mapStr = "{\"" + keyStr + "\":\"" + timeStr + "\"}"; + SchemaAndValue result = Values.parseString(mapStr); + assertEquals(Type.MAP, result.schema().type()); + Schema keySchema = result.schema().keySchema(); + Schema valueSchema = result.schema().valueSchema(); + assertEquals(Type.STRING, keySchema.type()); + assertEquals(Type.INT32, valueSchema.type()); + assertEquals(Time.LOGICAL_NAME, valueSchema.name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIME_FORMAT_PATTERN).parse(timeStr); + assertEquals(Collections.singletonMap(keyStr, expected), result.value()); + } + + @Test + public void shouldParseTimeStringAsTimeInMap() throws Exception { + String keyStr = "k1"; + String timeStr = "14:34:54.346Z"; + String mapStr = "{\"" + keyStr + "\":" + timeStr + "}"; + SchemaAndValue result = Values.parseString(mapStr); + assertEquals(Type.MAP, result.schema().type()); + Schema keySchema = result.schema().keySchema(); + Schema valueSchema = result.schema().valueSchema(); + assertEquals(Type.STRING, keySchema.type()); + assertEquals(Type.INT32, valueSchema.type()); + assertEquals(Time.LOGICAL_NAME, valueSchema.name()); + java.util.Date expected = new SimpleDateFormat(Values.ISO_8601_TIME_FORMAT_PATTERN).parse(timeStr); + assertEquals(Collections.singletonMap(keyStr, expected), result.value()); + } + /** * This is technically invalid JSON, and we don't want to simply ignore the blank elements. */ From f53e431b77355b9dba30f28e304cba13fef7d93c Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Wed, 22 Jan 2020 10:43:08 -0600 Subject: [PATCH 5/5] Rebased on updated logic in Values class --- .../main/java/org/apache/kafka/connect/data/Values.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java index bd6fc217d541f..d99fbcabf86df 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/data/Values.java @@ -886,7 +886,7 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No throw new DataException("Map is missing terminating '}': " + parser.original()); } } catch (DataException e) { - LOG.debug("Unable to parse the value as a map or an array; reverting to string", e); + LOG.trace("Unable to parse the value as a map or an array; reverting to string", e); parser.rewindTo(startPosition); } @@ -953,8 +953,10 @@ protected static SchemaAndValue parse(Parser parser, boolean embedded) throws No // can't parse as a number } } - // At this point, the only thing this can be is a string. Embedded strings were processed above, - // so this is not embedded and we can use the original string... + if (embedded) { + throw new DataException("Failed to parse embedded value"); + } + // At this point, the only thing this non-embedded value can be is a string. return new SchemaAndValue(Schema.STRING_SCHEMA, parser.original()); }