diff --git a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/message/CollectRep.java b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/message/CollectRep.java index 03d2821c5d4..0bac03f823a 100644 --- a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/message/CollectRep.java +++ b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/message/CollectRep.java @@ -285,8 +285,9 @@ public List getValues() { Row row = iterator.next(); ValueRow valueRow = ValueRow.newBuilder() .setColumns(fieldNames.stream() - .map(fieldName -> new String(((VarCharVector) - table.getVector(fieldName)).get(row.getRowNumber()))) + .map(fieldName -> new String( + ((VarCharVector) table.getVector(fieldName)).get(row.getRowNumber()), + StandardCharsets.UTF_8)) .collect(Collectors.toList())) .build(); values.add(valueRow); @@ -443,8 +444,8 @@ public MetricsData build() { fieldIndex < row.getColumnsList().size()) { String value = row.getColumns(fieldIndex); if (value != null) { - // Check byte array size, Arrow buffer size is 32768 bytes byte[] bytes = value.getBytes(StandardCharsets.UTF_8); + // setSafe grows the variable-width data buffer beyond its initial allocation. vector.setSafe(rowIndex, bytes); } } @@ -464,7 +465,7 @@ public MetricsData build() { throw e; } } - + public long getId() { return Long.parseLong(metadata.getOrDefault(MetricDataConstants.ID, "0")); } diff --git a/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/entity/message/CollectRepTest.java b/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/entity/message/CollectRepTest.java index 2cfcb51a6b6..f7202d363de 100644 --- a/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/entity/message/CollectRepTest.java +++ b/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/entity/message/CollectRepTest.java @@ -19,10 +19,16 @@ package org.apache.hertzbeat.common.entity.message; +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.stream.Stream; import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.MethodSource; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.params.provider.Arguments.arguments; /** * Test case for {@link CollectRep} @@ -43,4 +49,41 @@ void testFieldEquals(String name1, String name2, boolean result) { assertEquals(field1.equals(field2), result); } + @ParameterizedTest(name = "{0}") + @MethodSource("largeMetricValues") + void preservesArrowStringValue(String description, String value) { + CollectRep.Field field = CollectRep.Field.newBuilder() + .setName("payload") + .setType(1) + .build(); + CollectRep.ValueRow row = new CollectRep.ValueRow(List.of(value)); + + try (CollectRep.MetricsData metricsData = CollectRep.MetricsData.newBuilder() + .addField(field) + .addValueRow(row) + .build()) { + String storedValue = metricsData.getValues().getFirst().getColumns(0); + assertEquals( + value.getBytes(StandardCharsets.UTF_8).length, + storedValue.getBytes(StandardCharsets.UTF_8).length); + assertEquals(value, storedValue); + } + } + + private static Stream largeMetricValues() { + String ideograph = "\u4e2d"; + String emoji = new String(Character.toChars(0x1F600)); + return Stream.of( + arguments("ASCII before previous boundary", "a".repeat(32_699)), + arguments("ASCII at previous boundary", "a".repeat(32_700)), + arguments("ASCII after previous boundary", "a".repeat(32_701)), + arguments("large ASCII value", "a".repeat(100_000)), + arguments("multibyte before previous boundary", ideograph.repeat(10_899)), + arguments("multibyte at previous boundary", ideograph.repeat(10_900)), + arguments("multibyte after previous boundary", ideograph.repeat(10_901)), + arguments("emoji before previous boundary", emoji.repeat(8_174)), + arguments("emoji at previous boundary", emoji.repeat(8_175)), + arguments("emoji after previous boundary", emoji.repeat(8_176))); + } + } diff --git a/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/serialize/KafkaMetricsDataSerializerTest.java b/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/serialize/KafkaMetricsDataSerializerTest.java index dbb7807201c..a021712005c 100644 --- a/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/serialize/KafkaMetricsDataSerializerTest.java +++ b/hertzbeat-common-core/src/test/java/org/apache/hertzbeat/common/serialize/KafkaMetricsDataSerializerTest.java @@ -23,6 +23,8 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; import java.nio.channels.Channels; +import java.nio.charset.StandardCharsets; +import java.util.List; import java.util.Map; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.arrow.vector.ipc.ArrowStreamWriter; @@ -102,6 +104,29 @@ void testSerializeWithHeaders() { assertArrayEquals(expectedBytes, bytes); } + @Test + void preservesLargeMetricValueThroughArrowIpc() { + String value = "a".repeat(100_000); + CollectRep.Field field = CollectRep.Field.newBuilder() + .setName("payload") + .setType(1) + .build(); + CollectRep.MetricsData source = CollectRep.MetricsData.newBuilder() + .addField(field) + .addValueRow(new CollectRep.ValueRow(List.of(value))) + .build(); + + byte[] bytes = serializer.serialize("topic", source); + + KafkaMetricsDataDeserializer deserializer = new KafkaMetricsDataDeserializer(); + try (CollectRep.MetricsData restored = deserializer.deserialize("topic", bytes)) { + assertArrayEquals( + value.getBytes(StandardCharsets.UTF_8), + restored.getValues().getFirst().getColumns(0) + .getBytes(StandardCharsets.UTF_8)); + } + } + @Test void testClose() {