From c6063948e2029715412aedf5e629179433edeb31 Mon Sep 17 00:00:00 2001 From: Beetle brank <120192315+beetle0915@users.noreply.github.com> Date: Mon, 24 Aug 2026 17:12:35 +0800 Subject: [PATCH 1/4] [common] Collect min/max statistics for CHAR --- .../fluss/record/LogRecordBatchStatisticsCollector.java | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsCollector.java b/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsCollector.java index c4123d3ba13..d446a22cf76 100644 --- a/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsCollector.java +++ b/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsCollector.java @@ -34,6 +34,8 @@ import java.util.Arrays; import java.util.function.Function; +import static org.apache.fluss.types.DataTypeChecks.getLength; + /** * Collector for {@link LogRecordBatchStatistics} that accumulates statistics during record batch * construction. Manages statistics data in memory arrays and supports schema-aware statistics @@ -160,6 +162,11 @@ private void updateMinMax(int statsIndex, int schemaIndex, InternalRow row) { double doubleValue = row.getDouble(schemaIndex); updateMinMaxInternal(statsIndex, doubleValue, Double::compare); break; + case CHAR: + BinaryString charValue = row.getChar(schemaIndex, getLength(fieldType)); + updateMinMaxInternal( + statsIndex, charValue, BinaryString::compareTo, BinaryString::copy); + break; case STRING: BinaryString stringValue = row.getString(schemaIndex); updateMinMaxInternal( From 0f6bd6aedb6ecf57af17dd993a5229b317e64ed0 Mon Sep 17 00:00:00 2001 From: Beetle brank <120192315+beetle0915@users.noreply.github.com> Date: Mon, 24 Aug 2026 17:13:25 +0800 Subject: [PATCH 2/4] [common] Serialize CHAR statistics --- .../apache/fluss/record/LogRecordBatchStatisticsWriter.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsWriter.java b/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsWriter.java index c6faad9d87c..ef7f4354d91 100644 --- a/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsWriter.java +++ b/fluss-common/src/main/java/org/apache/fluss/record/LogRecordBatchStatisticsWriter.java @@ -31,6 +31,7 @@ import java.io.IOException; import static org.apache.fluss.record.LogRecordBatchFormat.STATISTICS_VERSION; +import static org.apache.fluss.types.DataTypeChecks.getLength; /** * A high-performance writer for LogRecordBatch statistics that efficiently serializes statistical @@ -205,6 +206,10 @@ private AlignedRow convertToAlignedRow(InternalRow row) { case DOUBLE: reusableRowWriter.writeDouble(i, row.getDouble(i)); break; + case CHAR: + int charLength = getLength(fieldType); + reusableRowWriter.writeChar(i, row.getChar(i, charLength), charLength); + break; case STRING: reusableRowWriter.writeString(i, row.getString(i)); break; From 2dd128dd261dba8ce447675394428a63f9de2429 Mon Sep 17 00:00:00 2001 From: Beetle brank <120192315+beetle0915@users.noreply.github.com> Date: Mon, 24 Aug 2026 17:15:54 +0800 Subject: [PATCH 3/4] [common][test] Test CHAR statistics round trip --- ...LogRecordBatchStatisticsCollectorTest.java | 28 +++++++++++++++++++ 1 file changed, 28 insertions(+) diff --git a/fluss-common/src/test/java/org/apache/fluss/record/LogRecordBatchStatisticsCollectorTest.java b/fluss-common/src/test/java/org/apache/fluss/record/LogRecordBatchStatisticsCollectorTest.java index 1b8506fee44..4b01d95e6d5 100644 --- a/fluss-common/src/test/java/org/apache/fluss/record/LogRecordBatchStatisticsCollectorTest.java +++ b/fluss-common/src/test/java/org/apache/fluss/record/LogRecordBatchStatisticsCollectorTest.java @@ -169,6 +169,34 @@ void testProcessNullValues() throws IOException { assertThat(statistics.getMaxValues().getInt(0)).isEqualTo(3); } + @Test + void testCharStatisticsSerializationRoundTrip() throws IOException { + RowType charRowType = DataTypes.ROW(new DataField("char_val", DataTypes.CHAR(3))); + LogRecordBatchStatisticsCollector charCollector = + new LogRecordBatchStatisticsCollector(charRowType, new int[] {0}); + + List charData = + Arrays.asList( + new Object[] {"cat"}, + new Object[] {"a"}, + new Object[] {null}, + new Object[] {"dog"}); + for (Object[] data : charData) { + charCollector.processRow(DataTestUtils.row(data)); + } + + MemorySegment segment = MemorySegment.allocateHeapMemory(1024); + charCollector.writeStatistics(new MemorySegmentOutputView(segment)); + + DefaultLogRecordBatchStatistics statistics = + LogRecordBatchStatisticsParser.parseStatistics(segment, 0, charRowType, 1); + + assertThat(statistics.getMinValues().getChar(0, 3)).isEqualTo(BinaryString.fromString("a")); + assertThat(statistics.getMaxValues().getChar(0, 3)) + .isEqualTo(BinaryString.fromString("dog")); + assertThat(statistics.getNullCounts()[0]).isEqualTo(1); + } + @Test void testPartialStatsIndexMapping() throws IOException { RowType fullRowType = From 54019b9895004fc0c1d64068ec99e847937136c6 Mon Sep 17 00:00:00 2001 From: Beetle brank <120192315+beetle0915@users.noreply.github.com> Date: Mon, 24 Aug 2026 17:18:08 +0800 Subject: [PATCH 4/4] [server][test] Test CHAR statistics pruning --- .../server/log/RecordBatchFilterTest.java | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/fluss-server/src/test/java/org/apache/fluss/server/log/RecordBatchFilterTest.java b/fluss-server/src/test/java/org/apache/fluss/server/log/RecordBatchFilterTest.java index b33da18e372..8d5b5e9f864 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/log/RecordBatchFilterTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/log/RecordBatchFilterTest.java @@ -125,6 +125,27 @@ void testSchemaAwareStatisticsWithNullStats() { assertThat(result).isTrue(); } + @Test + void testCharStatisticsCanPruneRecordBatch() throws IOException { + RowType rowType = RowType.of(DataTypes.CHAR(3)); + LogRecordBatchStatisticsCollector collector = + new LogRecordBatchStatisticsCollector(rowType, new int[] {0}); + collector.processRow(GenericRow.of(BinaryString.fromString("cat"))); + collector.processRow(GenericRow.of(BinaryString.fromString("ant"))); + collector.processRow(GenericRow.of(BinaryString.fromString("dog"))); + + MemorySegmentOutputView outputView = new MemorySegmentOutputView(1024); + collector.writeStatistics(outputView); + DefaultLogRecordBatchStatistics statistics = + LogRecordBatchStatisticsParser.parseStatistics( + outputView.getMemorySegment(), 0, rowType, TEST_SCHEMA_ID); + + Predicate predicate = + new PredicateBuilder(rowType).greaterThan(0, BinaryString.fromString("zoo")); + + assertThat(mayMatch(predicate, TEST_SCHEMA_ID, null, 3L, statistics)).isFalse(); + } + // ===================================================== // Schema Evolution Tests // =====================================================