From 07a4bd3fe89a843d05bfa4a4e10d22cae706e0bb Mon Sep 17 00:00:00 2001 From: martin Date: Wed, 16 Oct 2024 11:07:13 +0100 Subject: [PATCH 1/7] KAFKA-17792 header parsing times out processing and using large quantities of memory if the string looks like a number --- .../main/java/org/apache/kafka/connect/data/Values.java | 6 ++++++ .../java/org/apache/kafka/connect/data/ValuesTest.java | 7 +++++++ 2 files changed, 13 insertions(+) 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 cd332d2a8b3b3..5fe1b026c539b 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 @@ -71,6 +71,8 @@ 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 BigDecimal TOO_BIG = new BigDecimal("1e1000000"); + private static BigDecimal TOO_SMALL = new BigDecimal("1e-1000000"); private static final Pattern TWO_BACKSLASHES = Pattern.compile("\\\\"); @@ -1033,6 +1035,10 @@ private static SchemaAndValue parseAsNumber(String token) { } private static SchemaAndValue parseAsExactDecimal(BigDecimal decimal) { + BigDecimal abs = decimal.abs(); + if (abs.compareTo(TOO_BIG) > 0 || (abs.compareTo(TOO_SMALL) < 0 && BigDecimal.ZERO.compareTo(abs) != 0)) { + throw new NumberFormatException("outside efficient parsing range"); + } BigDecimal ceil = decimal.setScale(0, RoundingMode.CEILING); BigDecimal floor = decimal.setScale(0, RoundingMode.FLOOR); if (ceil.equals(floor)) { 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 e81b1f8a06ec1..3e630bda0a788 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 @@ -1180,7 +1180,14 @@ public void shouldParseFractionalPartsAsIntegerWhenNoFractionalPart() { assertEquals(new SchemaAndValue(Schema.INT32_SCHEMA, 66000), Values.parseString("66000.0")); assertEquals(new SchemaAndValue(Schema.FLOAT32_SCHEMA, 66000.0008f), Values.parseString("66000.0008")); } + @Test + public void avoidCpuAndMemoryIssuesConvertingExtremeBigDecimals() { + String PARSING_BIG = "1e+100000000"; // new BigDecimal().setScale(0, RoundingMode.FLOOR) takes around two minutes and uses 3GB; + assertEquals(new SchemaAndValue(Schema.STRING_SCHEMA, PARSING_BIG), Values.parseString(PARSING_BIG)); + String PARSING_SMALL = "1e-100000000"; + assertEquals(new SchemaAndValue(Schema.STRING_SCHEMA, PARSING_SMALL), Values.parseString(PARSING_SMALL)); + } protected void assertParsed(String input) { assertParsed(input, input); } From 91639cf81330b40a1e75107009f4baf767529064 Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 08:24:14 +0000 Subject: [PATCH 2/7] fix code and tests --- .../org/apache/kafka/connect/data/Values.java | 34 +++++++++---------- .../apache/kafka/connect/data/ValuesTest.java | 25 +++++++------- 2 files changed, 30 insertions(+), 29 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 5fe1b026c539b..ac12d5be6f77c 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 @@ -16,17 +16,9 @@ */ package org.apache.kafka.connect.data; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.errors.DataException; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import java.io.Serializable; import java.math.BigDecimal; import java.math.BigInteger; -import java.math.RoundingMode; import java.nio.ByteBuffer; import java.text.CharacterIterator; import java.text.DateFormat; @@ -44,6 +36,12 @@ import java.util.TimeZone; import java.util.regex.Pattern; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.errors.DataException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * Utility for converting from one Connect value to a different form. This is useful when the caller expects a value of a particular type * but is uncertain whether the actual value is one that isn't directly that type but can be converted into that type. @@ -71,8 +69,6 @@ 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 BigDecimal TOO_BIG = new BigDecimal("1e1000000"); - private static BigDecimal TOO_SMALL = new BigDecimal("1e-1000000"); private static final Pattern TWO_BACKSLASHES = Pattern.compile("\\\\"); @@ -1014,6 +1010,10 @@ private SchemaAndValue parseMultipleTokensAsTemporal(String token) { return parseAsTemporal(token); } + private static boolean isWholeNumber(BigDecimal bd) { + return bd.signum() == 0 || bd.scale() <= 0 || bd.stripTrailingZeros().scale() <= 0; + } + private static SchemaAndValue parseAsNumber(String token) { // Try to parse as a number ... BigDecimal decimal = new BigDecimal(token); @@ -1034,16 +1034,16 @@ private static SchemaAndValue parseAsNumber(String token) { } } + private static final BigDecimal BIGGER_THAN_LONG = new BigDecimal("1e19"); + private static SchemaAndValue parseAsExactDecimal(BigDecimal decimal) { BigDecimal abs = decimal.abs(); - if (abs.compareTo(TOO_BIG) > 0 || (abs.compareTo(TOO_SMALL) < 0 && BigDecimal.ZERO.compareTo(abs) != 0)) { - throw new NumberFormatException("outside efficient parsing range"); + if (abs.compareTo(BIGGER_THAN_LONG) > 0 || (abs.compareTo(BigDecimal.ONE) < 0 && abs.compareTo(BigDecimal.ZERO) != 0)) { + return null; } - BigDecimal ceil = decimal.setScale(0, RoundingMode.CEILING); - BigDecimal floor = decimal.setScale(0, RoundingMode.FLOOR); - if (ceil.equals(floor)) { - BigInteger num = ceil.toBigIntegerExact(); - if (ceil.precision() >= 19 && (num.compareTo(LONG_MIN) < 0 || num.compareTo(LONG_MAX) > 0)) { + if (isWholeNumber(decimal)) { + BigInteger num = decimal.toBigIntegerExact(); + if (num.compareTo(LONG_MIN) < 0 || num.compareTo(LONG_MAX) > 0) { return null; } long integral = num.longValue(); 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 3e630bda0a788..96814c7977437 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 @@ -16,14 +16,6 @@ */ package org.apache.kafka.connect.data; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.data.Values.Parser; -import org.apache.kafka.connect.errors.DataException; - -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.Timeout; - import java.math.BigDecimal; import java.math.BigInteger; import java.nio.ByteBuffer; @@ -45,6 +37,10 @@ import java.util.List; import java.util.Map; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.data.Values.Parser; +import org.apache.kafka.connect.errors.DataException; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -53,6 +49,8 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; public class ValuesTest { @@ -1182,12 +1180,15 @@ public void shouldParseFractionalPartsAsIntegerWhenNoFractionalPart() { } @Test public void avoidCpuAndMemoryIssuesConvertingExtremeBigDecimals() { - String PARSING_BIG = "1e+100000000"; // new BigDecimal().setScale(0, RoundingMode.FLOOR) takes around two minutes and uses 3GB; - assertEquals(new SchemaAndValue(Schema.STRING_SCHEMA, PARSING_BIG), Values.parseString(PARSING_BIG)); + String parsingBig = "1e+100000000"; // new BigDecimal().setScale(0, RoundingMode.FLOOR) takes around two minutes and uses 3GB; + BigDecimal valueBig = new BigDecimal(parsingBig); + assertEquals(new SchemaAndValue(Decimal.schema(-100000000), valueBig), Values.parseString(parsingBig), "parsing number that's too big"); - String PARSING_SMALL = "1e-100000000"; - assertEquals(new SchemaAndValue(Schema.STRING_SCHEMA, PARSING_SMALL), Values.parseString(PARSING_SMALL)); + String parsingSmall = "1e-100000000"; + BigDecimal valueSmall = new BigDecimal(parsingSmall); + assertEquals(new SchemaAndValue(Schema.FLOAT32_SCHEMA, (float) valueSmall.doubleValue()), Values.parseString(parsingSmall), "parsing number that's too big, strictly this should return a bigdecimal"); } + protected void assertParsed(String input) { assertParsed(input, input); } From bb00f700bec0a519ec83f59b60dabed480eb9edf Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 10:30:46 +0000 Subject: [PATCH 3/7] KAFKA-17792 header parsing times out processing and using large quantities of memory if the string looks like a number --- .../src/test/java/org/apache/kafka/connect/data/ValuesTest.java | 1 + 1 file changed, 1 insertion(+) 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 96814c7977437..555a609cf9470 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 @@ -1178,6 +1178,7 @@ public void shouldParseFractionalPartsAsIntegerWhenNoFractionalPart() { assertEquals(new SchemaAndValue(Schema.INT32_SCHEMA, 66000), Values.parseString("66000.0")); assertEquals(new SchemaAndValue(Schema.FLOAT32_SCHEMA, 66000.0008f), Values.parseString("66000.0008")); } + @Test public void avoidCpuAndMemoryIssuesConvertingExtremeBigDecimals() { String parsingBig = "1e+100000000"; // new BigDecimal().setScale(0, RoundingMode.FLOOR) takes around two minutes and uses 3GB; From 37fae21ec516f8abb80f129dc4e597d124eb4615 Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 10:36:51 +0000 Subject: [PATCH 4/7] still convert to int --- .../main/java/org/apache/kafka/connect/data/Values.java | 7 ++++--- 1 file changed, 4 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 ac12d5be6f77c..db4511de59d7c 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 @@ -1010,9 +1010,6 @@ private SchemaAndValue parseMultipleTokensAsTemporal(String token) { return parseAsTemporal(token); } - private static boolean isWholeNumber(BigDecimal bd) { - return bd.signum() == 0 || bd.scale() <= 0 || bd.stripTrailingZeros().scale() <= 0; - } private static SchemaAndValue parseAsNumber(String token) { // Try to parse as a number ... @@ -1034,6 +1031,10 @@ private static SchemaAndValue parseAsNumber(String token) { } } + private static boolean isWholeNumber(BigDecimal bd) { + return bd.signum() == 0 || bd.scale() <= 0 || bd.stripTrailingZeros().scale() <= 0; + } + private static final BigDecimal BIGGER_THAN_LONG = new BigDecimal("1e19"); private static SchemaAndValue parseAsExactDecimal(BigDecimal decimal) { From 52251558a935b1f312fbcab9b9ea3f4e0a657745 Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 08:31:20 +0000 Subject: [PATCH 5/7] spotless --- .../java/org/apache/kafka/connect/data/Values.java | 13 +++++++------ .../org/apache/kafka/connect/data/ValuesTest.java | 14 ++++++++------ 2 files changed, 15 insertions(+), 12 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 db4511de59d7c..f23f1f88a750c 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 @@ -16,6 +16,13 @@ */ package org.apache.kafka.connect.data; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.errors.DataException; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.io.Serializable; import java.math.BigDecimal; import java.math.BigInteger; @@ -36,12 +43,6 @@ import java.util.TimeZone; import java.util.regex.Pattern; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.errors.DataException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - /** * Utility for converting from one Connect value to a different form. This is useful when the caller expects a value of a particular type * but is uncertain whether the actual value is one that isn't directly that type but can be converted into that type. 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 555a609cf9470..ac6eef6fa6800 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 @@ -16,6 +16,14 @@ */ package org.apache.kafka.connect.data; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.data.Values.Parser; +import org.apache.kafka.connect.errors.DataException; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + import java.math.BigDecimal; import java.math.BigInteger; import java.nio.ByteBuffer; @@ -37,10 +45,6 @@ import java.util.List; import java.util.Map; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.data.Values.Parser; -import org.apache.kafka.connect.errors.DataException; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -49,8 +53,6 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.Timeout; public class ValuesTest { From b407e5813e362525be5bef8da97e14dcaf6a97fe Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 10:18:57 +0000 Subject: [PATCH 6/7] add missing new line --- .../org/apache/kafka/connect/data/ValuesTest.java | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) 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 ac6eef6fa6800..555a609cf9470 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 @@ -16,14 +16,6 @@ */ package org.apache.kafka.connect.data; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.data.Values.Parser; -import org.apache.kafka.connect.errors.DataException; - -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.Timeout; - import java.math.BigDecimal; import java.math.BigInteger; import java.nio.ByteBuffer; @@ -45,6 +37,10 @@ import java.util.List; import java.util.Map; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.data.Values.Parser; +import org.apache.kafka.connect.errors.DataException; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -53,6 +49,8 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; public class ValuesTest { From a4d376b9a354ccf4f802a845e30aeda53c528596 Mon Sep 17 00:00:00 2001 From: martin Date: Fri, 24 Jan 2025 10:42:23 +0000 Subject: [PATCH 7/7] spotless --- .../org/apache/kafka/connect/data/ValuesTest.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) 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 555a609cf9470..ac6eef6fa6800 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 @@ -16,6 +16,14 @@ */ package org.apache.kafka.connect.data; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema.Type; +import org.apache.kafka.connect.data.Values.Parser; +import org.apache.kafka.connect.errors.DataException; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + import java.math.BigDecimal; import java.math.BigInteger; import java.nio.ByteBuffer; @@ -37,10 +45,6 @@ import java.util.List; import java.util.Map; -import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.data.Schema.Type; -import org.apache.kafka.connect.data.Values.Parser; -import org.apache.kafka.connect.errors.DataException; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; @@ -49,8 +53,6 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.Timeout; public class ValuesTest {