From 25f56e38e587ce50fa58628ae57a2288ed87b673 Mon Sep 17 00:00:00 2001 From: Logic Date: Thu, 30 Jul 2026 16:50:29 +0800 Subject: [PATCH 1/2] maintenance: bound metric string values --- .../common/entity/message/CollectRep.java | 20 +++++++++++++++-- .../common/entity/message/CollectRepTest.java | 22 +++++++++++++++++++ 2 files changed, 40 insertions(+), 2 deletions(-) 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..ef9c6ec3be0 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 @@ -21,7 +21,10 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.CharBuffer; import java.nio.channels.Channels; +import java.nio.charset.CodingErrorAction; import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.HashMap; @@ -48,6 +51,8 @@ @SuppressWarnings("all") @Slf4j public final class CollectRep { + private static final int MAX_ARROW_VALUE_BYTES = 32_700; + private CollectRep() {} /** @@ -443,8 +448,7 @@ 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); + byte[] bytes = encodeMetricValue(value); vector.setSafe(rowIndex, bytes); } } @@ -464,6 +468,18 @@ public MetricsData build() { throw e; } } + + private static byte[] encodeMetricValue(String value) { + ByteBuffer buffer = ByteBuffer.allocate(MAX_ARROW_VALUE_BYTES); + StandardCharsets.UTF_8.newEncoder() + .onMalformedInput(CodingErrorAction.REPLACE) + .onUnmappableCharacter(CodingErrorAction.REPLACE) + .encode(CharBuffer.wrap(value), buffer, true); + buffer.flip(); + byte[] encoded = new byte[buffer.remaining()]; + buffer.get(encoded); + return encoded; + } 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..1c29bb496d5 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,14 @@ package org.apache.hertzbeat.common.entity.message; +import java.nio.charset.StandardCharsets; +import java.util.List; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; /** * Test case for {@link CollectRep} @@ -43,4 +47,22 @@ void testFieldEquals(String name1, String name2, boolean result) { assertEquals(field1.equals(field2), result); } + @Test + void boundsArrowStringValueSize() { + String oversizedValue = "a".repeat(100_000); + CollectRep.Field field = CollectRep.Field.newBuilder() + .setName("payload") + .setType(1) + .build(); + CollectRep.ValueRow row = new CollectRep.ValueRow(List.of(oversizedValue)); + + try (CollectRep.MetricsData metricsData = CollectRep.MetricsData.newBuilder() + .addField(field) + .addValueRow(row) + .build()) { + String storedValue = metricsData.getValues().getFirst().getColumns(0); + assertTrue(storedValue.getBytes(StandardCharsets.UTF_8).length <= 32_700); + } + } + } From 5dcdbfbf8e83444f24ee4320b88706e553c2ce35 Mon Sep 17 00:00:00 2001 From: Logic Date: Thu, 30 Jul 2026 21:37:32 +0800 Subject: [PATCH 2/2] preserve arrow metric value fidelity --- .../common/entity/message/CollectRep.java | 25 +++---------- .../common/entity/message/CollectRepTest.java | 35 +++++++++++++++---- .../KafkaMetricsDataSerializerTest.java | 25 +++++++++++++ 3 files changed, 58 insertions(+), 27 deletions(-) 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 ef9c6ec3be0..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 @@ -21,10 +21,7 @@ import java.io.ByteArrayOutputStream; import java.io.IOException; -import java.nio.ByteBuffer; -import java.nio.CharBuffer; import java.nio.channels.Channels; -import java.nio.charset.CodingErrorAction; import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.HashMap; @@ -51,8 +48,6 @@ @SuppressWarnings("all") @Slf4j public final class CollectRep { - private static final int MAX_ARROW_VALUE_BYTES = 32_700; - private CollectRep() {} /** @@ -290,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); @@ -448,7 +444,8 @@ public MetricsData build() { fieldIndex < row.getColumnsList().size()) { String value = row.getColumns(fieldIndex); if (value != null) { - byte[] bytes = encodeMetricValue(value); + byte[] bytes = value.getBytes(StandardCharsets.UTF_8); + // setSafe grows the variable-width data buffer beyond its initial allocation. vector.setSafe(rowIndex, bytes); } } @@ -469,18 +466,6 @@ public MetricsData build() { } } - private static byte[] encodeMetricValue(String value) { - ByteBuffer buffer = ByteBuffer.allocate(MAX_ARROW_VALUE_BYTES); - StandardCharsets.UTF_8.newEncoder() - .onMalformedInput(CodingErrorAction.REPLACE) - .onUnmappableCharacter(CodingErrorAction.REPLACE) - .encode(CharBuffer.wrap(value), buffer, true); - buffer.flip(); - byte[] encoded = new byte[buffer.remaining()]; - buffer.get(encoded); - return encoded; - } - 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 1c29bb496d5..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 @@ -21,12 +21,14 @@ import java.nio.charset.StandardCharsets; import java.util.List; -import org.junit.jupiter.api.Test; +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.api.Assertions.assertTrue; +import static org.junit.jupiter.params.provider.Arguments.arguments; /** * Test case for {@link CollectRep} @@ -47,22 +49,41 @@ void testFieldEquals(String name1, String name2, boolean result) { assertEquals(field1.equals(field2), result); } - @Test - void boundsArrowStringValueSize() { - String oversizedValue = "a".repeat(100_000); + @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(oversizedValue)); + 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); - assertTrue(storedValue.getBytes(StandardCharsets.UTF_8).length <= 32_700); + 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() {