maintenance: preserve metric string values (#4287)

Co-authored-by: Duansg <siguoduan@gmail.com>
This commit is contained in:
Logic
2026-08-06 14:46:59 +08:00
committed by GitHub
co-authored by Duansg
parent bf3bb4fa45
commit 7eb4b1cd44
3 changed files with 73 additions and 4 deletions
@@ -285,8 +285,9 @@ public final class CollectRep {
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 final class CollectRep {
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 final class CollectRep {
throw e;
}
}
public long getId() {
return Long.parseLong(metadata.getOrDefault(MetricDataConstants.ID, "0"));
}
@@ -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 @@ public class CollectRepTest {
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<Arguments> 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)));
}
}
@@ -23,6 +23,8 @@ import static org.junit.jupiter.api.Assertions.assertNull;
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 @@ class KafkaMetricsDataSerializerTest {
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() {