Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -285,8 +285,9 @@ public List<ValueRow> 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);
Expand Down Expand Up @@ -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);
}
}
Expand All @@ -464,7 +465,7 @@ public MetricsData build() {
throw e;
}
}

public long getId() {
return Long.parseLong(metadata.getOrDefault(MetricDataConstants.ID, "0"));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand All @@ -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<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)));
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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() {

Expand Down
Loading