From 62adff4e3f83f65758adeda6dd92e19479f3724b Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sat, 18 Jul 2026 19:51:31 +0800 Subject: [PATCH 01/10] fix(iceberg): honor disabled write metrics Jira: DORIS-27023 Filter backend-collected file statistics using the Iceberg table metrics configuration before writing manifest metadata. --- .../iceberg/helper/IcebergWriterHelper.java | 31 ++++++++++++-- .../helper/IcebergWriterHelperTest.java | 42 +++++++++++++++++++ 2 files changed, 69 insertions(+), 4 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java index b67a5911b64384..99b99ac9a46939 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java @@ -30,8 +30,12 @@ import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileMetadata; import org.apache.iceberg.Metrics; +import org.apache.iceberg.MetricsConfig; +import org.apache.iceberg.MetricsModes; +import org.apache.iceberg.MetricsUtil; import org.apache.iceberg.PartitionData; import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.io.WriteResult; @@ -70,7 +74,7 @@ public static WriteResult convertToWriterResult( long fileSize = commitData.getFileSize(); long recordCount = commitData.getRowCount(); CommonStatistics stat = new CommonStatistics(recordCount, DEFAULT_FILE_COUNT, fileSize); - Metrics metrics = buildDataFileMetrics(table, fileFormat, commitData); + Metrics metrics = buildDataFileMetrics(table, commitData); Optional partitionData = Optional.empty(); //get and check partitionValues when table is partitionedTable if (spec.isPartitioned()) { @@ -153,7 +157,7 @@ private static PartitionData convertToPartitionData( return partitionData; } - private static Metrics buildDataFileMetrics(Table table, FileFormat fileFormat, TIcebergCommitData commitData) { + private static Metrics buildDataFileMetrics(Table table, TIcebergCommitData commitData) { Map columnSizes = new HashMap<>(); Map valueCounts = new HashMap<>(); Map nullValueCounts = new HashMap<>(); @@ -178,8 +182,27 @@ private static Metrics buildDataFileMetrics(Table table, FileFormat fileFormat, } } - return new Metrics(commitData.getRowCount(), columnSizes, valueCounts, - nullValueCounts, null, lowerBounds, upperBounds); + MetricsConfig metricsConfig = MetricsConfig.forTable(table); + Schema schema = table.schema(); + // Physical file stats may contain every column, but manifest metrics must honor the table's metadata policy. + return new Metrics(commitData.getRowCount(), + filterDisabledMetrics(columnSizes, schema, metricsConfig), + filterDisabledMetrics(valueCounts, schema, metricsConfig), + filterDisabledMetrics(nullValueCounts, schema, metricsConfig), + null, + filterDisabledMetrics(lowerBounds, schema, metricsConfig), + filterDisabledMetrics(upperBounds, schema, metricsConfig)); + } + + private static Map filterDisabledMetrics( + Map metrics, Schema schema, MetricsConfig metricsConfig) { + Map filteredMetrics = new HashMap<>(); + metrics.forEach((fieldId, value) -> { + if (MetricsUtil.metricsMode(schema, metricsConfig, fieldId) != MetricsModes.None.get()) { + filteredMetrics.put(fieldId, value); + } + }); + return filteredMetrics; } /** diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index 77d318518f326e..aad7ff5e0047b1 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -18,19 +18,28 @@ package org.apache.doris.datasource.iceberg.helper; import org.apache.doris.thrift.TFileContent; +import org.apache.doris.thrift.TIcebergColumnStats; import org.apache.doris.thrift.TIcebergCommitData; +import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.Schema; +import org.apache.iceberg.SortOrder; +import org.apache.iceberg.Table; +import org.apache.iceberg.TableProperties; +import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; +import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.List; +import java.util.Map; /** * Test for IcebergWriterHelper DeleteFile conversion @@ -58,6 +67,39 @@ public void setUp() { } + @Test + public void testConvertToWriterResultRespectsNoneMetricsMode() { + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(schema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none")); + + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setColumnSizes(Map.of(2, 128L)); + columnStats.setValueCounts(Map.of(2, 10L)); + columnStats.setNullValueCounts(Map.of(2, 0L)); + columnStats.setLowerBounds(Map.of(2, ByteBuffer.wrap(new byte[] {0x01}))); + columnStats.setUpperBounds(Map.of(2, ByteBuffer.wrap(new byte[] {0x02}))); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/data.parquet"); + commitData.setRowCount(10); + commitData.setFileSize(1024); + commitData.setColumnStats(columnStats); + + WriteResult result = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)); + DataFile dataFile = result.dataFiles()[0]; + + Assertions.assertTrue(dataFile.columnSizes() == null || dataFile.columnSizes().isEmpty()); + Assertions.assertTrue(dataFile.valueCounts() == null || dataFile.valueCounts().isEmpty()); + Assertions.assertTrue(dataFile.nullValueCounts() == null || dataFile.nullValueCounts().isEmpty()); + Assertions.assertTrue(dataFile.lowerBounds() == null || dataFile.lowerBounds().isEmpty()); + Assertions.assertTrue(dataFile.upperBounds() == null || dataFile.upperBounds().isEmpty()); + } + @Test public void testConvertToDeleteFiles_EmptyList() { List commitDataList = new ArrayList<>(); From c95c9d1292ca8f40f76dcb6ba6a5c2e40fa7c0b4 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sat, 18 Jul 2026 21:22:30 +0800 Subject: [PATCH 02/10] fix(iceberg): handle absent metrics in stats Build the Iceberg metrics policy once per commit batch and return unknown statistics when required file metrics are absent.\n\nJira: DORIS-27023 --- .../iceberg/helper/IcebergWriterHelper.java | 9 ++-- .../doris/statistics/util/StatisticsUtil.java | 12 ++++- .../helper/IcebergWriterHelperTest.java | 26 ++++++++++ .../statistics/util/StatisticsUtilTest.java | 47 +++++++++++++++++++ 4 files changed, 88 insertions(+), 6 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java index 99b99ac9a46939..ee89407333b8c9 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java @@ -65,6 +65,8 @@ public static WriteResult convertToWriterResult( // Get table specification information PartitionSpec spec = table.spec(); FileFormat fileFormat = IcebergUtils.getFileFormat(table); + MetricsConfig metricsConfig = MetricsConfig.forTable(table); + Schema schema = table.schema(); for (TIcebergCommitData commitData : commitDataList) { //get the files path @@ -74,7 +76,7 @@ public static WriteResult convertToWriterResult( long fileSize = commitData.getFileSize(); long recordCount = commitData.getRowCount(); CommonStatistics stat = new CommonStatistics(recordCount, DEFAULT_FILE_COUNT, fileSize); - Metrics metrics = buildDataFileMetrics(table, commitData); + Metrics metrics = buildDataFileMetrics(commitData, schema, metricsConfig); Optional partitionData = Optional.empty(); //get and check partitionValues when table is partitionedTable if (spec.isPartitioned()) { @@ -157,7 +159,8 @@ private static PartitionData convertToPartitionData( return partitionData; } - private static Metrics buildDataFileMetrics(Table table, TIcebergCommitData commitData) { + private static Metrics buildDataFileMetrics( + TIcebergCommitData commitData, Schema schema, MetricsConfig metricsConfig) { Map columnSizes = new HashMap<>(); Map valueCounts = new HashMap<>(); Map nullValueCounts = new HashMap<>(); @@ -182,8 +185,6 @@ private static Metrics buildDataFileMetrics(Table table, TIcebergCommitData comm } } - MetricsConfig metricsConfig = MetricsConfig.forTable(table); - Schema schema = table.schema(); // Physical file stats may contain every column, but manifest metrics must honor the table's metadata policy. return new Metrics(commitData.getRowCount(), filterDisabledMetrics(columnSizes, schema, metricsConfig), diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java index 10b139b73f047f..7556b83a3a4c19 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java @@ -605,9 +605,17 @@ public static Optional getIcebergColumnStats(String colName, or try (CloseableIterable fileScanTasks = tableScan.planFiles()) { for (FileScanTask task : fileScanTasks) { int colId = getColId(task.spec(), colName); - totalDataSize += task.file().columnSizes().get(colId); + Map columnSizes = task.file().columnSizes(); + Map nullValueCounts = task.file().nullValueCounts(); + Long columnSize = columnSizes == null ? null : columnSizes.get(colId); + Long nullValueCount = nullValueCounts == null ? null : nullValueCounts.get(colId); + // Iceberg can omit maps or entries for mode=none; partial aggregation would fabricate zero stats. + if (columnSize == null || nullValueCount == null) { + return Optional.empty(); + } + totalDataSize += columnSize; totalDataCount += task.file().recordCount(); - totalNumNull += task.file().nullValueCounts().get(colId); + totalNumNull += nullValueCount; } } catch (IOException e) { LOG.warn("Error to close FileScanTask.", e); diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index aad7ff5e0047b1..da2adf225db892 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -100,6 +100,32 @@ public void testConvertToWriterResultRespectsNoneMetricsMode() { Assertions.assertTrue(dataFile.upperBounds() == null || dataFile.upperBounds().isEmpty()); } + @Test + public void testConvertToWriterResultBuildsMetricsPolicyOncePerBatch() { + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(schema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none")); + + TIcebergCommitData firstCommit = new TIcebergCommitData(); + firstCommit.setFilePath("/path/to/first.parquet"); + firstCommit.setRowCount(10); + firstCommit.setFileSize(1024); + + TIcebergCommitData secondCommit = new TIcebergCommitData(); + secondCommit.setFilePath("/path/to/second.parquet"); + secondCommit.setRowCount(20); + secondCommit.setFileSize(2048); + + IcebergWriterHelper.convertToWriterResult(table, List.of(firstCommit, secondCommit)); + + // One schema lookup is made by Iceberg's policy builder and one is captured for all files in the batch. + Mockito.verify(table, Mockito.times(2)).schema(); + } + @Test public void testConvertToDeleteFiles_EmptyList() { List commitDataList = new ArrayList<>(); diff --git a/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java b/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java index 8bccbd2929f665..033c605cf254c9 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java @@ -44,16 +44,28 @@ import org.apache.doris.datasource.iceberg.IcebergExternalDatabase; import org.apache.doris.datasource.iceberg.IcebergExternalTable; import org.apache.doris.datasource.iceberg.IcebergHadoopExternalCatalog; +import org.apache.doris.datasource.iceberg.helper.IcebergWriterHelper; import org.apache.doris.nereids.trees.expressions.literal.Literal; import org.apache.doris.qe.SessionVariable; import org.apache.doris.rpc.RpcException; import org.apache.doris.statistics.AnalysisManager; import org.apache.doris.statistics.ColStatsMeta; import org.apache.doris.statistics.TableStatsMeta; +import org.apache.doris.thrift.TIcebergColumnStats; +import org.apache.doris.thrift.TIcebergCommitData; import org.apache.doris.thrift.TStorageType; import com.google.common.collect.Maps; import org.apache.iceberg.CatalogProperties; +import org.apache.iceberg.DataFile; +import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; +import org.apache.iceberg.SortOrder; +import org.apache.iceberg.TableProperties; +import org.apache.iceberg.TableScan; +import org.apache.iceberg.io.CloseableIterable; +import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.mockito.MockedStatic; @@ -67,6 +79,41 @@ import java.util.Map; class StatisticsUtilTest { + @Test + void testGetIcebergColumnStatsReturnsEmptyForDisabledMetrics() { + Schema schema = new Schema(Types.NestedField.optional(1, "id", Types.IntegerType.get())); + PartitionSpec spec = PartitionSpec.builderFor(schema).build(); + org.apache.iceberg.Table table = Mockito.mock(org.apache.iceberg.Table.class); + Mockito.when(table.schema()).thenReturn(schema); + Mockito.when(table.spec()).thenReturn(spec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none")); + + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setColumnSizes(Map.of(1, 128L)); + columnStats.setValueCounts(Map.of(1, 10L)); + columnStats.setNullValueCounts(Map.of(1, 0L)); + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/data.parquet"); + commitData.setRowCount(10); + commitData.setFileSize(1024); + commitData.setColumnStats(columnStats); + DataFile dataFile = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]; + + TableScan tableScan = Mockito.mock(TableScan.class); + FileScanTask fileScanTask = Mockito.mock(FileScanTask.class); + Mockito.when(table.newScan()).thenReturn(tableScan); + Mockito.when(tableScan.includeColumnStats()).thenReturn(tableScan); + Mockito.when(tableScan.planFiles()) + .thenReturn(CloseableIterable.withNoopClose(List.of(fileScanTask))); + Mockito.when(fileScanTask.spec()).thenReturn(spec); + Mockito.when(fileScanTask.file()).thenReturn(dataFile); + + Assertions.assertTrue(StatisticsUtil.getIcebergColumnStats("id", table).isEmpty()); + } + @Test void testConvertToDouble() { try { From 42d2ca9dddb9634e85c28dc585939a64428a315a Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 05:46:10 +0800 Subject: [PATCH 03/10] fix(iceberg): preserve unknown stats on close failure Return unknown Iceberg column statistics when closing the scan fails, preventing partial or zero accumulators from escaping. Jira: DORIS-27023 --- .../doris/statistics/util/StatisticsUtil.java | 2 ++ .../statistics/util/StatisticsUtilTest.java | 23 +++++++++++++++++++ 2 files changed, 25 insertions(+) diff --git a/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java b/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java index 7556b83a3a4c19..d683824699917e 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java +++ b/fe/fe-core/src/main/java/org/apache/doris/statistics/util/StatisticsUtil.java @@ -619,6 +619,8 @@ public static Optional getIcebergColumnStats(String colName, or } } catch (IOException e) { LOG.warn("Error to close FileScanTask.", e); + // A failed close can cancel an in-flight empty return, so accumulated stats are not reliable. + return Optional.empty(); } ColumnStatisticBuilder columnStatisticBuilder = new ColumnStatisticBuilder(totalDataCount); columnStatisticBuilder.setMaxValue(Double.POSITIVE_INFINITY); diff --git a/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java b/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java index 033c605cf254c9..af9bd9d2ac378f 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/statistics/util/StatisticsUtilTest.java @@ -71,6 +71,7 @@ import org.mockito.MockedStatic; import org.mockito.Mockito; +import java.io.IOException; import java.time.LocalTime; import java.time.format.DateTimeFormatter; import java.util.ArrayList; @@ -114,6 +115,28 @@ void testGetIcebergColumnStatsReturnsEmptyForDisabledMetrics() { Assertions.assertTrue(StatisticsUtil.getIcebergColumnStats("id", table).isEmpty()); } + @Test + void testGetIcebergColumnStatsReturnsEmptyWhenCloseFails() { + Schema schema = new Schema(Types.NestedField.optional(1, "id", Types.IntegerType.get())); + PartitionSpec spec = PartitionSpec.builderFor(schema).build(); + org.apache.iceberg.Table table = Mockito.mock(org.apache.iceberg.Table.class); + TableScan tableScan = Mockito.mock(TableScan.class); + FileScanTask fileScanTask = Mockito.mock(FileScanTask.class); + DataFile dataFile = Mockito.mock(DataFile.class); + Mockito.when(table.newScan()).thenReturn(tableScan); + Mockito.when(tableScan.includeColumnStats()).thenReturn(tableScan); + Mockito.when(tableScan.planFiles()).thenReturn(CloseableIterable.combine( + List.of(fileScanTask), () -> { + throw new IOException("close failed"); + })); + Mockito.when(fileScanTask.spec()).thenReturn(spec); + Mockito.when(fileScanTask.file()).thenReturn(dataFile); + Mockito.when(dataFile.columnSizes()).thenReturn(Map.of()); + Mockito.when(dataFile.nullValueCounts()).thenReturn(Map.of()); + + Assertions.assertTrue(StatisticsUtil.getIcebergColumnStats("id", table).isEmpty()); + } + @Test void testConvertToDouble() { try { From a91df699cc3a0c226f2e521f52535cfa1c3bb053 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 09:16:38 +0800 Subject: [PATCH 04/10] fix(iceberg): apply metrics modes to bounds Omit bounds for counts mode and safely truncate string and binary bounds for truncate mode before writing Iceberg manifests. Jira: DORIS-27023 --- .../iceberg/helper/IcebergWriterHelper.java | 48 ++++++++++++- .../helper/IcebergWriterHelperTest.java | 72 +++++++++++++++++++ 2 files changed, 118 insertions(+), 2 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java index ee89407333b8c9..84f53c4eb3c888 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java @@ -39,7 +39,11 @@ import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.io.WriteResult; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; import org.apache.iceberg.types.Types; +import org.apache.iceberg.util.BinaryUtil; +import org.apache.iceberg.util.UnicodeUtil; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -191,8 +195,8 @@ private static Metrics buildDataFileMetrics( filterDisabledMetrics(valueCounts, schema, metricsConfig), filterDisabledMetrics(nullValueCounts, schema, metricsConfig), null, - filterDisabledMetrics(lowerBounds, schema, metricsConfig), - filterDisabledMetrics(upperBounds, schema, metricsConfig)); + filterBounds(lowerBounds, schema, metricsConfig, true), + filterBounds(upperBounds, schema, metricsConfig, false)); } private static Map filterDisabledMetrics( @@ -206,6 +210,46 @@ private static Map filterDisabledMetrics( return filteredMetrics; } + private static Map filterBounds( + Map bounds, Schema schema, MetricsConfig metricsConfig, boolean lowerBound) { + Map filteredBounds = new HashMap<>(); + bounds.forEach((fieldId, value) -> { + MetricsModes.MetricsMode mode = MetricsUtil.metricsMode(schema, metricsConfig, fieldId); + if (mode == MetricsModes.None.get() || mode == MetricsModes.Counts.get()) { + return; + } + + ByteBuffer filteredValue = value; + if (mode instanceof MetricsModes.Truncate) { + Type type = schema.findType(fieldId); + int length = ((MetricsModes.Truncate) mode).length(); + // Truncated upper bounds must round up so file pruning cannot exclude matching values. + filteredValue = truncateBound(type, value, length, lowerBound); + } + if (filteredValue != null) { + filteredBounds.put(fieldId, filteredValue); + } + }); + return filteredBounds; + } + + private static ByteBuffer truncateBound(Type type, ByteBuffer value, int length, boolean lowerBound) { + switch (type.typeId()) { + case STRING: + String stringValue = Conversions.fromByteBuffer(type, value).toString(); + String truncatedString = lowerBound + ? UnicodeUtil.truncateStringMin(stringValue, length) + : UnicodeUtil.truncateStringMax(stringValue, length); + return truncatedString == null ? null : Conversions.toByteBuffer(type, truncatedString); + case BINARY: + return lowerBound + ? BinaryUtil.truncateBinaryMin(value, length) + : BinaryUtil.truncateBinaryMax(value, length); + default: + return value; + } + } + /** * Convert TIcebergCommitData list to DeleteFile list for delete operations. * diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index da2adf225db892..a025881b8fdc4a 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -30,6 +30,7 @@ import org.apache.iceberg.Table; import org.apache.iceberg.TableProperties; import org.apache.iceberg.io.WriteResult; +import org.apache.iceberg.types.Conversions; import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; @@ -100,6 +101,77 @@ public void testConvertToWriterResultRespectsNoneMetricsMode() { Assertions.assertTrue(dataFile.upperBounds() == null || dataFile.upperBounds().isEmpty()); } + @Test + public void testConvertToWriterResultCountsModeOmitsBounds() { + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(schema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "counts")); + + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setColumnSizes(Map.of(2, 128L)); + columnStats.setValueCounts(Map.of(2, 10L)); + columnStats.setNullValueCounts(Map.of(2, 0L)); + columnStats.setLowerBounds(Map.of( + 2, Conversions.toByteBuffer(Types.StringType.get(), "abcdefgh"))); + columnStats.setUpperBounds(Map.of( + 2, Conversions.toByteBuffer(Types.StringType.get(), "ijklmnop"))); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/data.parquet"); + commitData.setRowCount(10); + commitData.setFileSize(1024); + commitData.setColumnStats(columnStats); + + DataFile dataFile = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]; + + Assertions.assertEquals(Map.of(2, 128L), dataFile.columnSizes()); + Assertions.assertEquals(Map.of(2, 10L), dataFile.valueCounts()); + Assertions.assertEquals(Map.of(2, 0L), dataFile.nullValueCounts()); + Assertions.assertTrue(dataFile.lowerBounds() == null || dataFile.lowerBounds().isEmpty()); + Assertions.assertTrue(dataFile.upperBounds() == null || dataFile.upperBounds().isEmpty()); + } + + @Test + public void testConvertToWriterResultTruncatesStringAndBinaryBounds() { + Schema boundsSchema = new Schema( + Types.NestedField.optional(1, "text", Types.StringType.get()), + Types.NestedField.optional(2, "payload", Types.BinaryType.get())); + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(boundsSchema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "truncate(3)")); + + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setLowerBounds(Map.of( + 1, Conversions.toByteBuffer(Types.StringType.get(), "abcdef"), + 2, ByteBuffer.wrap(new byte[] {1, 2, 3, 4}))); + columnStats.setUpperBounds(Map.of( + 1, Conversions.toByteBuffer(Types.StringType.get(), "uvwxyz"), + 2, ByteBuffer.wrap(new byte[] {1, 2, 3, 4}))); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/data.parquet"); + commitData.setRowCount(10); + commitData.setFileSize(1024); + commitData.setColumnStats(columnStats); + + DataFile dataFile = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]; + + Assertions.assertEquals("abc", Conversions.fromByteBuffer( + Types.StringType.get(), dataFile.lowerBounds().get(1)).toString()); + Assertions.assertEquals("uvx", Conversions.fromByteBuffer( + Types.StringType.get(), dataFile.upperBounds().get(1)).toString()); + Assertions.assertEquals(ByteBuffer.wrap(new byte[] {1, 2, 3}), dataFile.lowerBounds().get(2)); + Assertions.assertEquals(ByteBuffer.wrap(new byte[] {1, 2, 4}), dataFile.upperBounds().get(2)); + } + @Test public void testConvertToWriterResultBuildsMetricsPolicyOncePerBatch() { Table table = Mockito.mock(Table.class); From 1fe01c8a48c5ad2d040b4c0d6dddf3719cc1bc37 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 11:06:38 +0800 Subject: [PATCH 05/10] fix(iceberg): address write metrics review feedback --- be/src/exec/sink/viceberg_merge_sink.cpp | 3 + .../iceberg/viceberg_partition_writer.cpp | 10 +- .../iceberg/viceberg_partition_writer.h | 3 + .../iceberg/iceberg_partition_writer_test.cpp | 107 ++++++++++++++++++ .../datasource/iceberg/IcebergUtils.java | 12 ++ .../iceberg/helper/IcebergWriterHelper.java | 46 +++++++- .../doris/planner/IcebergMergeSink.java | 1 + .../doris/planner/IcebergTableSink.java | 1 + .../helper/IcebergWriterHelperTest.java | 72 ++++++++++++ .../doris/planner/IcebergMergeSinkTest.java | 32 ++++++ gensrc/thrift/DataSinks.thrift | 4 + 11 files changed, 285 insertions(+), 6 deletions(-) create mode 100644 be/test/exec/sink/writer/iceberg/iceberg_partition_writer_test.cpp diff --git a/be/src/exec/sink/viceberg_merge_sink.cpp b/be/src/exec/sink/viceberg_merge_sink.cpp index 7d6c84cf068dbc..e9f4d6bf5c655a 100644 --- a/be/src/exec/sink/viceberg_merge_sink.cpp +++ b/be/src/exec/sink/viceberg_merge_sink.cpp @@ -254,6 +254,9 @@ Status VIcebergMergeSink::_build_inner_sinks() { if (merge_sink.__isset.broker_addresses) { table_sink.__set_broker_addresses(merge_sink.broker_addresses); } + if (merge_sink.__isset.collect_column_stats) { + table_sink.__set_collect_column_stats(merge_sink.collect_column_stats); + } _table_sink.__set_type(TDataSinkType::ICEBERG_TABLE_SINK); _table_sink.__set_iceberg_table_sink(table_sink); diff --git a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp index 5f28c41c24d61b..5958165ad0395c 100644 --- a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp +++ b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.cpp @@ -47,7 +47,11 @@ VIcebergPartitionWriter::VIcebergPartitionWriter( _file_name_index(file_name_index), _file_format_type(file_format_type), _compress_type(compress_type), - _hadoop_conf(hadoop_conf) {} + _hadoop_conf(hadoop_conf) { + if (t_sink.iceberg_table_sink.__isset.collect_column_stats) { + _collect_column_stats = t_sink.iceberg_table_sink.collect_column_stats; + } +} Status VIcebergPartitionWriter::open(RuntimeState* state, RuntimeProfile* profile, const RowDescriptor* row_desc) { @@ -156,6 +160,10 @@ Status VIcebergPartitionWriter::_build_iceberg_commit_data(TIcebergCommitData* c commit_data->__set_file_size(_file_format_transformer->written_len()); commit_data->__set_file_content(TFileContent::DATA); commit_data->__set_partition_values(_partition_values); + // ORC collection reopens the file, so honor the FE policy before any footer work. + if (!_collect_column_stats) { + return Status::OK(); + } if (_file_format_type == TFileFormatType::FORMAT_PARQUET) { TIcebergColumnStats column_stats; RETURN_IF_ERROR(static_cast(_file_format_transformer.get()) diff --git a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h index 97b4dd3efdfac0..b0839a82ed3796 100644 --- a/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h +++ b/be/src/exec/sink/writer/iceberg/viceberg_partition_writer.h @@ -67,6 +67,8 @@ class VIcebergPartitionWriter : public IPartitionWriterBase { inline size_t written_len() const override { return _file_format_transformer->written_len(); } private: + friend class VIcebergPartitionWriterTest; + std::string _get_target_file_name(); Status _build_iceberg_commit_data(TIcebergCommitData* commit_data); @@ -91,6 +93,7 @@ class VIcebergPartitionWriter : public IPartitionWriterBase { TFileFormatType::type _file_format_type; TFileCompressType::type _compress_type; const std::map& _hadoop_conf; + bool _collect_column_stats = true; std::shared_ptr _fs = nullptr; diff --git a/be/test/exec/sink/writer/iceberg/iceberg_partition_writer_test.cpp b/be/test/exec/sink/writer/iceberg/iceberg_partition_writer_test.cpp new file mode 100644 index 00000000000000..d453177cf25044 --- /dev/null +++ b/be/test/exec/sink/writer/iceberg/iceberg_partition_writer_test.cpp @@ -0,0 +1,107 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include + +#include + +#include "exec/sink/writer/iceberg/viceberg_partition_writer.h" + +namespace doris { + +namespace { + +class FakeFileFormatTransformer final : public VFileFormatTransformer { +public: + explicit FakeFileFormatTransformer(const VExprContextSPtrs& output_exprs) + : VFileFormatTransformer(nullptr, output_exprs, false) {} + + Status open() override { return Status::OK(); } + Status write(const Block&) override { return Status::OK(); } + Status close() override { return Status::OK(); } + int64_t written_len() override { return 64; } +}; + +TDataSink make_table_sink(std::optional collect_column_stats) { + TIcebergTableSink iceberg_sink; + if (collect_column_stats.has_value()) { + iceberg_sink.__set_collect_column_stats(*collect_column_stats); + } + TDataSink sink; + sink.__set_type(TDataSinkType::ICEBERG_TABLE_SINK); + sink.__set_iceberg_table_sink(iceberg_sink); + return sink; +} + +} // namespace + +class VIcebergPartitionWriterTest : public testing::Test { +protected: + static std::unique_ptr make_writer( + const TDataSink& sink, const VExprContextSPtrs& output_exprs, + const iceberg::Schema& schema, const std::string* schema_json, + const std::map& hadoop_conf) { + IPartitionWriterBase::WriteInfo write_info; + write_info.file_type = TFileType::FILE_LOCAL; + return std::make_unique( + sink, std::vector {}, output_exprs, schema, schema_json, + std::vector {}, std::move(write_info), "data", 0, + TFileFormatType::FORMAT_ORC, TFileCompressType::ZLIB, hadoop_conf); + } + + static void install_fake_transformer(VIcebergPartitionWriter* writer, + const VExprContextSPtrs& output_exprs) { + writer->_file_format_transformer = + std::make_unique(output_exprs); + } + + static Status build_commit_data(VIcebergPartitionWriter* writer, + TIcebergCommitData* commit_data) { + return writer->_build_iceberg_commit_data(commit_data); + } + + static bool collect_column_stats(const VIcebergPartitionWriter& writer) { + return writer._collect_column_stats; + } +}; + +TEST_F(VIcebergPartitionWriterTest, OrcSkipsFooterCollectionWhenMetricsAreDisabled) { + VExprContextSPtrs output_exprs; + iceberg::Schema schema(std::vector {}); + std::string schema_json; + std::map hadoop_conf; + auto writer = + make_writer(make_table_sink(false), output_exprs, schema, &schema_json, hadoop_conf); + install_fake_transformer(writer.get(), output_exprs); + + TIcebergCommitData commit_data; + ASSERT_TRUE(build_commit_data(writer.get(), &commit_data).ok()); + EXPECT_FALSE(commit_data.__isset.column_stats); +} + +TEST_F(VIcebergPartitionWriterTest, MissingPolicyKeepsCollectionEnabledForRollingUpgrade) { + VExprContextSPtrs output_exprs; + iceberg::Schema schema(std::vector {}); + std::string schema_json; + std::map hadoop_conf; + auto writer = make_writer(make_table_sink(std::nullopt), output_exprs, schema, &schema_json, + hadoop_conf); + + EXPECT_TRUE(collect_column_stats(*writer)); +} + +} // namespace doris diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java index b6f1d22046d76e..562b1e9564479a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java @@ -87,6 +87,9 @@ import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.MetadataTableType; import org.apache.iceberg.MetadataTableUtils; +import org.apache.iceberg.MetricsConfig; +import org.apache.iceberg.MetricsModes; +import org.apache.iceberg.MetricsUtil; import org.apache.iceberg.PartitionData; import org.apache.iceberg.PartitionField; import org.apache.iceberg.PartitionSpec; @@ -1961,6 +1964,15 @@ public static Schema appendRowLineageFieldsForV3(Schema schema) { MetadataColumns.ROW_ID, MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER)); } + public static boolean shouldCollectColumnStats(Table table) { + Schema schema = table.schema(); + MetricsConfig metricsConfig = MetricsConfig.forTable(table); + return TypeUtil.indexById(schema.asStruct()).values().stream() + .filter(field -> field.type().isPrimitiveType()) + .anyMatch(field -> MetricsUtil.metricsMode(schema, metricsConfig, field.fieldId()) + != MetricsModes.None.get()); + } + public static int getFormatVersion(Table table) { int formatVersion = 2; // default format version : 2 if (table instanceof BaseTable) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java index 84f53c4eb3c888..094cbef038dd03 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java @@ -41,6 +41,7 @@ import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.types.Conversions; import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.TypeUtil; import org.apache.iceberg.types.Types; import org.apache.iceberg.util.BinaryUtil; import org.apache.iceberg.util.UnicodeUtil; @@ -71,6 +72,10 @@ public static WriteResult convertToWriterResult( FileFormat fileFormat = IcebergUtils.getFileFormat(table); MetricsConfig metricsConfig = MetricsConfig.forTable(table); Schema schema = table.schema(); + if (IcebergUtils.getFormatVersion(table) >= IcebergUtils.ICEBERG_ROW_LINEAGE_MIN_VERSION) { + // Rewrite and merge writers emit v3 lineage columns that are absent from the table schema. + schema = IcebergUtils.appendRowLineageFieldsForV3(schema); + } for (TIcebergCommitData commitData : commitDataList) { //get the files path @@ -165,6 +170,7 @@ private static PartitionData convertToPartitionData( private static Metrics buildDataFileMetrics( TIcebergCommitData commitData, Schema schema, MetricsConfig metricsConfig) { + Map fieldParents = TypeUtil.indexParents(schema.asStruct()); Map columnSizes = new HashMap<>(); Map valueCounts = new HashMap<>(); Map nullValueCounts = new HashMap<>(); @@ -192,11 +198,11 @@ private static Metrics buildDataFileMetrics( // Physical file stats may contain every column, but manifest metrics must honor the table's metadata policy. return new Metrics(commitData.getRowCount(), filterDisabledMetrics(columnSizes, schema, metricsConfig), - filterDisabledMetrics(valueCounts, schema, metricsConfig), - filterDisabledMetrics(nullValueCounts, schema, metricsConfig), + filterLogicalMetrics(valueCounts, schema, metricsConfig, fieldParents), + filterLogicalMetrics(nullValueCounts, schema, metricsConfig, fieldParents), null, - filterBounds(lowerBounds, schema, metricsConfig, true), - filterBounds(upperBounds, schema, metricsConfig, false)); + filterBounds(lowerBounds, schema, metricsConfig, fieldParents, true), + filterBounds(upperBounds, schema, metricsConfig, fieldParents, false)); } private static Map filterDisabledMetrics( @@ -210,10 +216,28 @@ private static Map filterDisabledMetrics( return filteredMetrics; } + private static Map filterLogicalMetrics( + Map metrics, Schema schema, MetricsConfig metricsConfig, + Map fieldParents) { + Map filteredMetrics = new HashMap<>(); + metrics.forEach((fieldId, value) -> { + // Definition-level values below list/map do not represent logical element counts. + if (!isInRepeatedField(fieldId, schema, fieldParents) + && MetricsUtil.metricsMode(schema, metricsConfig, fieldId) != MetricsModes.None.get()) { + filteredMetrics.put(fieldId, value); + } + }); + return filteredMetrics; + } + private static Map filterBounds( - Map bounds, Schema schema, MetricsConfig metricsConfig, boolean lowerBound) { + Map bounds, Schema schema, MetricsConfig metricsConfig, + Map fieldParents, boolean lowerBound) { Map filteredBounds = new HashMap<>(); bounds.forEach((fieldId, value) -> { + if (isInRepeatedField(fieldId, schema, fieldParents)) { + return; + } MetricsModes.MetricsMode mode = MetricsUtil.metricsMode(schema, metricsConfig, fieldId); if (mode == MetricsModes.None.get() || mode == MetricsModes.Counts.get()) { return; @@ -233,6 +257,18 @@ private static Map filterBounds( return filteredBounds; } + private static boolean isInRepeatedField( + int fieldId, Schema schema, Map fieldParents) { + Integer parentId = fieldId; + while ((parentId = fieldParents.get(parentId)) != null) { + Types.NestedField parent = schema.findField(parentId); + if (parent != null && !parent.type().isStructType()) { + return true; + } + } + return false; + } + private static ByteBuffer truncateBound(Type type, ByteBuffer value, int length, boolean lowerBound) { switch (type.typeId()) { case STRING: diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java index 4af4ba17e18578..75e204d22be37e 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java @@ -131,6 +131,7 @@ public void bindDataSink(Optional insertCtx) } tSink.setFormatVersion(formatVersion); tSink.setSchemaJson(SchemaParser.toJson(schema)); + tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable)); // partition spec if (icebergTable.spec().isPartitioned()) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java index 0f3b1bb24d26bc..3cf169769c8bab 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java @@ -135,6 +135,7 @@ public void bindDataSink(Optional insertCtx) schema = IcebergUtils.appendRowLineageFieldsForV3(schema); } tSink.setSchemaJson(SchemaParser.toJson(schema)); + tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable)); // partition spec if (icebergTable.spec().isPartitioned()) { diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index a025881b8fdc4a..465cc0fed9884a 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -24,6 +24,7 @@ import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; +import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.PartitionSpec; import org.apache.iceberg.Schema; import org.apache.iceberg.SortOrder; @@ -198,6 +199,77 @@ public void testConvertToWriterResultBuildsMetricsPolicyOncePerBatch() { Mockito.verify(table, Mockito.times(2)).schema(); } + @Test + public void testConvertToWriterResultHandlesV3RowLineageMetrics() { + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(schema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.FORMAT_VERSION, "3", + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "truncate(16)")); + + int rowId = MetadataColumns.ROW_ID.fieldId(); + ByteBuffer rowIdBound = Conversions.toByteBuffer(MetadataColumns.ROW_ID.type(), 7L); + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setLowerBounds(Map.of(rowId, rowIdBound)); + columnStats.setUpperBounds(Map.of(rowId, rowIdBound)); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/v3-data.parquet"); + commitData.setRowCount(1); + commitData.setFileSize(128); + commitData.setColumnStats(columnStats); + + DataFile dataFile = Assertions.assertDoesNotThrow( + () -> IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]); + Assertions.assertEquals(rowIdBound, dataFile.lowerBounds().get(rowId)); + Assertions.assertEquals(rowIdBound, dataFile.upperBounds().get(rowId)); + } + + @Test + public void testConvertToWriterResultSuppressesLogicalMetricsBelowRepeatedFields() { + Schema repeatedSchema = new Schema( + Types.NestedField.optional(1, "items", + Types.ListType.ofOptional(2, Types.IntegerType.get())), + Types.NestedField.optional(3, "attributes", + Types.MapType.ofOptional(4, 5, Types.StringType.get(), Types.StringType.get())), + Types.NestedField.optional(6, "top_level", Types.IntegerType.get())); + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(repeatedSchema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "parquet", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "full")); + + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setColumnSizes(Map.of(2, 20L, 4, 40L, 5, 50L, 6, 60L)); + columnStats.setValueCounts(Map.of(2, 2L, 4, 4L, 5, 5L, 6, 6L)); + columnStats.setNullValueCounts(Map.of(2, 0L, 4, 0L, 5, 0L, 6, 0L)); + columnStats.setLowerBounds(Map.of( + 2, Conversions.toByteBuffer(Types.IntegerType.get(), 2), + 4, Conversions.toByteBuffer(Types.StringType.get(), "key"), + 5, Conversions.toByteBuffer(Types.StringType.get(), "value"), + 6, Conversions.toByteBuffer(Types.IntegerType.get(), 6))); + columnStats.setUpperBounds(columnStats.getLowerBounds()); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/repeated.parquet"); + commitData.setRowCount(6); + commitData.setFileSize(1024); + commitData.setColumnStats(columnStats); + + DataFile dataFile = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]; + + Assertions.assertEquals(Map.of(2, 20L, 4, 40L, 5, 50L, 6, 60L), dataFile.columnSizes()); + Assertions.assertEquals(Map.of(6, 6L), dataFile.valueCounts()); + Assertions.assertEquals(Map.of(6, 0L), dataFile.nullValueCounts()); + Assertions.assertEquals(Map.of(6, columnStats.getLowerBounds().get(6)), dataFile.lowerBounds()); + Assertions.assertEquals(Map.of(6, columnStats.getUpperBounds().get(6)), dataFile.upperBounds()); + } + @Test public void testConvertToDeleteFiles_EmptyList() { List commitDataList = new ArrayList<>(); diff --git a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java index 23dcb4403ce731..8938ee131fa84d 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java @@ -74,6 +74,32 @@ public void testBindDataSinkSkipsRewritableDeleteFileSetsAndRowLineageSchemaForV IcebergUtils.ICEBERG_LAST_UPDATED_SEQUENCE_NUMBER_COL)); } + @Test + public void testBindDataSinkDisablesColumnStatsWhenAllMetricsAreNone() throws Exception { + IcebergMergeSink sink = new IcebergMergeSink(mockIcebergExternalTable(2, Map.of( + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none")), new DeleteCommandContext()); + + sink.bindDataSink(Optional.empty()); + + TIcebergMergeSink thriftSink = sink.tDataSink.getIcebergMergeSink(); + Assertions.assertTrue(thriftSink.isSetCollectColumnStats()); + Assertions.assertFalse(thriftSink.isCollectColumnStats()); + } + + @Test + public void testBindDataSinkKeepsColumnStatsForMetricsOverride() throws Exception { + IcebergMergeSink sink = new IcebergMergeSink(mockIcebergExternalTable(2, Map.of( + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none", + TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX + "id", "counts")), + new DeleteCommandContext()); + + sink.bindDataSink(Optional.empty()); + + TIcebergMergeSink thriftSink = sink.tDataSink.getIcebergMergeSink(); + Assertions.assertTrue(thriftSink.isSetCollectColumnStats()); + Assertions.assertTrue(thriftSink.isCollectColumnStats()); + } + private static TIcebergRewritableDeleteFileSet buildDeleteFileSet() { TIcebergDeleteFileDesc deleteFileDesc = new TIcebergDeleteFileDesc(); deleteFileDesc.setPath("file:///tmp/delete.puffin"); @@ -84,6 +110,11 @@ private static TIcebergRewritableDeleteFileSet buildDeleteFileSet() { } private static IcebergExternalTable mockIcebergExternalTable(int formatVersion) { + return mockIcebergExternalTable(formatVersion, Collections.emptyMap()); + } + + private static IcebergExternalTable mockIcebergExternalTable( + int formatVersion, Map metricsProperties) { Schema schema = new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get())); PartitionSpec spec = PartitionSpec.unpartitioned(); Map properties = new HashMap<>(); @@ -91,6 +122,7 @@ private static IcebergExternalTable mockIcebergExternalTable(int formatVersion) properties.put(TableProperties.DEFAULT_FILE_FORMAT, "parquet"); properties.put(TableProperties.PARQUET_COMPRESSION, "snappy"); properties.put(TableProperties.WRITE_DATA_LOCATION, "file:///tmp/iceberg_tbl/data"); + properties.putAll(metricsProperties); Table icebergTable = Mockito.mock(Table.class); Mockito.when(icebergTable.properties()).thenReturn(properties); diff --git a/gensrc/thrift/DataSinks.thrift b/gensrc/thrift/DataSinks.thrift index 3a45208ea9bf11..18e152f447df9c 100644 --- a/gensrc/thrift/DataSinks.thrift +++ b/gensrc/thrift/DataSinks.thrift @@ -487,6 +487,8 @@ struct TIcebergTableSink { 15: optional map static_partition_values; 16: optional PlanNodes.TSortInfo sort_info; 17: optional TIcebergWriteType write_type = TIcebergWriteType.INSERT; + // Unset keeps collection enabled for rolling upgrades with older FEs. + 18: optional bool collect_column_stats; } struct TIcebergRewritableDeleteFileSet { @@ -533,6 +535,8 @@ struct TIcebergMergeSink { 11: optional map hadoop_config 12: optional Types.TFileType file_type 13: optional list broker_addresses; + // Unset keeps collection enabled for rolling upgrades with older FEs. + 14: optional bool collect_column_stats; // delete side (position delete only) 20: optional TFileContent delete_type From f05f437b42dfffd4802a41609fa6f548fa752893 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 14:50:24 +0800 Subject: [PATCH 06/10] fix(iceberg): detect v3 transaction table format --- .../datasource/iceberg/IcebergUtils.java | 7 +++-- .../helper/IcebergWriterHelperTest.java | 29 ++++++++++++------- 2 files changed, 23 insertions(+), 13 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java index 562b1e9564479a..56701346209771 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java @@ -79,10 +79,10 @@ import com.google.gson.reflect.TypeToken; import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.exception.ExceptionUtils; -import org.apache.iceberg.BaseTable; import org.apache.iceberg.CatalogProperties; import org.apache.iceberg.FileFormat; import org.apache.iceberg.FileScanTask; +import org.apache.iceberg.HasTableOperations; import org.apache.iceberg.ManifestFile; import org.apache.iceberg.MetadataColumns; import org.apache.iceberg.MetadataTableType; @@ -1975,8 +1975,9 @@ public static boolean shouldCollectColumnStats(Table table) { public static int getFormatVersion(Table table) { int formatVersion = 2; // default format version : 2 - if (table instanceof BaseTable) { - formatVersion = ((BaseTable) table).operations().current().formatVersion(); + if (table instanceof HasTableOperations) { + // TransactionTable exposes the real format version through operations, not table properties. + formatVersion = ((HasTableOperations) table).operations().current().formatVersion(); } else if (table != null && table.properties() != null) { String version = table.properties().get(TableProperties.FORMAT_VERSION); if (version != null) { diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index 465cc0fed9884a..1769abdcc4c7f6 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -21,6 +21,7 @@ import org.apache.doris.thrift.TIcebergColumnStats; import org.apache.doris.thrift.TIcebergCommitData; +import org.apache.hadoop.conf.Configuration; import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; @@ -30,15 +31,18 @@ import org.apache.iceberg.SortOrder; import org.apache.iceberg.Table; import org.apache.iceberg.TableProperties; +import org.apache.iceberg.hadoop.HadoopTables; import org.apache.iceberg.io.WriteResult; import org.apache.iceberg.types.Conversions; import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; import org.mockito.Mockito; import java.nio.ByteBuffer; +import java.nio.file.Path; import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -200,21 +204,23 @@ public void testConvertToWriterResultBuildsMetricsPolicyOncePerBatch() { } @Test - public void testConvertToWriterResultHandlesV3RowLineageMetrics() { - Table table = Mockito.mock(Table.class); - Mockito.when(table.schema()).thenReturn(schema); - Mockito.when(table.spec()).thenReturn(unpartitionedSpec); - Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); - Mockito.when(table.properties()).thenReturn(Map.of( + public void testConvertToWriterResultHandlesV3TransactionTableLineageMetrics(@TempDir Path tempDir) { + HadoopTables tables = new HadoopTables(new Configuration()); + Table baseTable = tables.create(schema, unpartitionedSpec, SortOrder.unsorted(), Map.of( TableProperties.FORMAT_VERSION, "3", TableProperties.DEFAULT_FILE_FORMAT, "parquet", - TableProperties.DEFAULT_WRITE_METRICS_MODE, "truncate(16)")); + TableProperties.DEFAULT_WRITE_METRICS_MODE, "truncate(16)"), + tempDir.resolve("table").toUri().toString()); + Table transactionTable = baseTable.newTransaction().table(); int rowId = MetadataColumns.ROW_ID.fieldId(); + int sequenceNumberId = MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.fieldId(); ByteBuffer rowIdBound = Conversions.toByteBuffer(MetadataColumns.ROW_ID.type(), 7L); + ByteBuffer sequenceNumberBound = Conversions.toByteBuffer( + MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER.type(), 3L); TIcebergColumnStats columnStats = new TIcebergColumnStats(); - columnStats.setLowerBounds(Map.of(rowId, rowIdBound)); - columnStats.setUpperBounds(Map.of(rowId, rowIdBound)); + columnStats.setLowerBounds(Map.of(rowId, rowIdBound, sequenceNumberId, sequenceNumberBound)); + columnStats.setUpperBounds(Map.of(rowId, rowIdBound, sequenceNumberId, sequenceNumberBound)); TIcebergCommitData commitData = new TIcebergCommitData(); commitData.setFilePath("/path/to/v3-data.parquet"); @@ -223,9 +229,12 @@ public void testConvertToWriterResultHandlesV3RowLineageMetrics() { commitData.setColumnStats(columnStats); DataFile dataFile = Assertions.assertDoesNotThrow( - () -> IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]); + () -> IcebergWriterHelper.convertToWriterResult( + transactionTable, List.of(commitData)).dataFiles()[0]); Assertions.assertEquals(rowIdBound, dataFile.lowerBounds().get(rowId)); Assertions.assertEquals(rowIdBound, dataFile.upperBounds().get(rowId)); + Assertions.assertEquals(sequenceNumberBound, dataFile.lowerBounds().get(sequenceNumberId)); + Assertions.assertEquals(sequenceNumberBound, dataFile.upperBounds().get(sequenceNumberId)); } @Test From 53a14e38c9155eefaa1207531f3c7b1d278fe39c Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 17:10:14 +0800 Subject: [PATCH 07/10] fix(iceberg): evaluate metrics for writer schema --- .../datasource/iceberg/IcebergUtils.java | 13 +++++-- .../doris/planner/IcebergMergeSink.java | 2 +- .../doris/planner/IcebergTableSink.java | 2 +- .../doris/planner/IcebergMergeSinkTest.java | 38 ++++++++++++++++++- 4 files changed, 48 insertions(+), 7 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java index 56701346209771..c9705b8492a8a7 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/IcebergUtils.java @@ -1964,12 +1964,17 @@ public static Schema appendRowLineageFieldsForV3(Schema schema) { MetadataColumns.ROW_ID, MetadataColumns.LAST_UPDATED_SEQUENCE_NUMBER)); } - public static boolean shouldCollectColumnStats(Table table) { - Schema schema = table.schema(); + public static boolean shouldCollectColumnStats(Table table, Schema writerSchema) { MetricsConfig metricsConfig = MetricsConfig.forTable(table); - return TypeUtil.indexById(schema.asStruct()).values().stream() + if (getFileFormat(table) == FileFormat.ORC) { + // Match the footer collectors: ORC reports top-level collection counts, while Parquet reports leaf fields. + return writerSchema.columns().stream() + .anyMatch(field -> MetricsUtil.metricsMode(writerSchema, metricsConfig, field.fieldId()) + != MetricsModes.None.get()); + } + return TypeUtil.indexById(writerSchema.asStruct()).values().stream() .filter(field -> field.type().isPrimitiveType()) - .anyMatch(field -> MetricsUtil.metricsMode(schema, metricsConfig, field.fieldId()) + .anyMatch(field -> MetricsUtil.metricsMode(writerSchema, metricsConfig, field.fieldId()) != MetricsModes.None.get()); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java index 75e204d22be37e..9775f81c172fae 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergMergeSink.java @@ -131,7 +131,7 @@ public void bindDataSink(Optional insertCtx) } tSink.setFormatVersion(formatVersion); tSink.setSchemaJson(SchemaParser.toJson(schema)); - tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable)); + tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable, schema)); // partition spec if (icebergTable.spec().isPartitioned()) { diff --git a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java index 3cf169769c8bab..b7d3da47cb4d45 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java +++ b/fe/fe-core/src/main/java/org/apache/doris/planner/IcebergTableSink.java @@ -135,7 +135,7 @@ public void bindDataSink(Optional insertCtx) schema = IcebergUtils.appendRowLineageFieldsForV3(schema); } tSink.setSchemaJson(SchemaParser.toJson(schema)); - tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable)); + tSink.setCollectColumnStats(IcebergUtils.shouldCollectColumnStats(icebergTable, schema)); // partition spec if (icebergTable.spec().isPartitioned()) { diff --git a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java index 8938ee131fa84d..df9b2ae80c26f2 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/planner/IcebergMergeSinkTest.java @@ -100,6 +100,37 @@ public void testBindDataSinkKeepsColumnStatsForMetricsOverride() throws Exceptio Assertions.assertTrue(thriftSink.isCollectColumnStats()); } + @Test + public void testBindDataSinkKeepsColumnStatsForV3LineageFields() throws Exception { + IcebergMergeSink sink = new IcebergMergeSink(mockIcebergExternalTable(3, Map.of( + TableProperties.DEFAULT_WRITE_METRICS_MODE, "counts", + TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX + "id", "none")), + new DeleteCommandContext()); + + sink.bindDataSink(Optional.empty()); + + TIcebergMergeSink thriftSink = sink.tDataSink.getIcebergMergeSink(); + Assertions.assertTrue(thriftSink.isSetCollectColumnStats()); + Assertions.assertTrue(thriftSink.isCollectColumnStats()); + } + + @Test + public void testBindDataSinkKeepsColumnStatsForOrcTopLevelComplexField() throws Exception { + Schema schema = new Schema(Types.NestedField.optional(1, "items", + Types.ListType.ofOptional(2, Types.IntegerType.get()))); + IcebergMergeSink sink = new IcebergMergeSink(mockIcebergExternalTable(2, schema, Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "orc", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "none", + TableProperties.METRICS_MODE_COLUMN_CONF_PREFIX + "items", "counts")), + new DeleteCommandContext()); + + sink.bindDataSink(Optional.empty()); + + TIcebergMergeSink thriftSink = sink.tDataSink.getIcebergMergeSink(); + Assertions.assertTrue(thriftSink.isSetCollectColumnStats()); + Assertions.assertTrue(thriftSink.isCollectColumnStats()); + } + private static TIcebergRewritableDeleteFileSet buildDeleteFileSet() { TIcebergDeleteFileDesc deleteFileDesc = new TIcebergDeleteFileDesc(); deleteFileDesc.setPath("file:///tmp/delete.puffin"); @@ -115,7 +146,12 @@ private static IcebergExternalTable mockIcebergExternalTable(int formatVersion) private static IcebergExternalTable mockIcebergExternalTable( int formatVersion, Map metricsProperties) { - Schema schema = new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get())); + return mockIcebergExternalTable(formatVersion, + new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get())), metricsProperties); + } + + private static IcebergExternalTable mockIcebergExternalTable( + int formatVersion, Schema schema, Map metricsProperties) { PartitionSpec spec = PartitionSpec.unpartitioned(); Map properties = new HashMap<>(); properties.put(TableProperties.FORMAT_VERSION, String.valueOf(formatVersion)); From 471470f5375d9483cd9b7034008fcedc8289f50b Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 19:04:47 +0800 Subject: [PATCH 08/10] fix(iceberg): map ORC stats by column id --- .../format/transformer/vorc_transformer.cpp | 7 +- .../transformer/vorc_transformer_test.cpp | 107 ++++++++++++++++++ 2 files changed, 113 insertions(+), 1 deletion(-) create mode 100644 be/test/format/transformer/vorc_transformer_test.cpp diff --git a/be/src/format/transformer/vorc_transformer.cpp b/be/src/format/transformer/vorc_transformer.cpp index 7dfa9fb4c64d3b..2004ce21d5a574 100644 --- a/be/src/format/transformer/vorc_transformer.cpp +++ b/be/src/format/transformer/vorc_transformer.cpp @@ -385,8 +385,13 @@ Status VOrcTransformer::collect_file_statistics_after_close(TIcebergColumnStats* const iceberg::StructType& root_struct = _iceberg_schema->root_struct(); const auto& nested_fields = root_struct.fields(); + const orc::Type& orc_root_type = reader->getType(); for (uint32_t i = 0; i < nested_fields.size(); i++) { - uint32_t orc_col_id = i + 1; // skip root struct + if (i >= orc_root_type.getSubtypeCount()) { + continue; + } + // ORC IDs are depth-first, so top-level fields after a complex field are not i + 1. + uint64_t orc_col_id = orc_root_type.getSubtype(i)->getColumnId(); if (orc_col_id >= file_stats->getNumberOfColumns()) { continue; } diff --git a/be/test/format/transformer/vorc_transformer_test.cpp b/be/test/format/transformer/vorc_transformer_test.cpp new file mode 100644 index 00000000000000..4ea14766356d15 --- /dev/null +++ b/be/test/format/transformer/vorc_transformer_test.cpp @@ -0,0 +1,107 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include "format/transformer/vorc_transformer.h" + +#include + +#include "core/block/block.h" +#include "core/column/column_string.h" +#include "core/column/column_struct.h" +#include "core/column/column_vector.h" +#include "core/data_type/data_type_number.h" +#include "core/data_type/data_type_string.h" +#include "core/data_type/data_type_struct.h" +#include "format/table/iceberg/schema_parser.h" +#include "io/fs/local_file_system.h" +#include "runtime/runtime_state.h" +#include "testutil/mock/mock_slot_ref.h" +#include "util/uid_util.h" + +namespace doris { + +class VOrcTransformerTest : public testing::Test { +protected: + void SetUp() override { + _file_path = "./vorc_transformer_" + UniqueId::gen_uid().to_string() + ".orc"; + _fs = io::global_local_filesystem(); + } + + void TearDown() override { static_cast(_fs->delete_file(_file_path)); } + + std::string _file_path; + std::shared_ptr _fs; +}; + +TEST_F(VOrcTransformerTest, CollectsBoundsForTopLevelFieldAfterStruct) { + auto int_type = std::make_shared(); + auto struct_type = std::make_shared(DataTypes {int_type}, Strings {"a"}); + auto string_type = std::make_shared(); + VExprContextSPtrs output_exprs = + MockSlotRef::create_mock_contexts(DataTypes {struct_type, string_type}); + + const std::string schema_json = R"({ + "type": "struct", + "fields": [ + { + "id": 1, + "name": "s", + "required": true, + "type": { + "type": "struct", + "fields": [ + {"id": 2, "name": "a", "required": true, "type": "int"} + ] + } + }, + {"id": 3, "name": "b", "required": true, "type": "string"} + ] + })"; + std::unique_ptr schema = iceberg::SchemaParser::from_json(schema_json); + + io::FileWriterPtr file_writer; + ASSERT_TRUE(_fs->create_file(_file_path, &file_writer).ok()); + RuntimeState state; + VOrcTransformer transformer(&state, file_writer.get(), output_exprs, "", {"s", "b"}, false, + TFileCompressType::PLAIN, schema.get(), _fs); + ASSERT_TRUE(transformer.open().ok()); + + auto nested_column = ColumnInt32::create(); + nested_column->insert_value(-1); + Columns struct_columns; + struct_columns.emplace_back(std::move(nested_column)); + auto struct_column = ColumnStruct::create(std::move(struct_columns)); + auto string_column = ColumnString::create(); + string_column->insert_data("hello", 5); + + Block block; + block.insert(ColumnWithTypeAndName(std::move(struct_column), struct_type, "s")); + block.insert(ColumnWithTypeAndName(std::move(string_column), string_type, "b")); + ASSERT_TRUE(transformer.write(block).ok()); + ASSERT_TRUE(transformer.close().ok()); + + TIcebergColumnStats stats; + ASSERT_TRUE(transformer.collect_file_statistics_after_close(&stats).ok()); + ASSERT_TRUE(stats.__isset.lower_bounds); + ASSERT_TRUE(stats.__isset.upper_bounds); + ASSERT_EQ(1, stats.lower_bounds.count(3)); + ASSERT_EQ(1, stats.upper_bounds.count(3)); + EXPECT_EQ("hello", stats.lower_bounds.at(3)); + EXPECT_EQ("hello", stats.upper_bounds.at(3)); +} + +} // namespace doris From d3abdfa0f08329c30b411c09646b293acb37f02a Mon Sep 17 00:00:00 2001 From: Gabriel Date: Sun, 19 Jul 2026 20:40:37 +0800 Subject: [PATCH 09/10] fix(iceberg): narrow ORC statistics column id safely --- be/src/format/transformer/vorc_transformer.cpp | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/be/src/format/transformer/vorc_transformer.cpp b/be/src/format/transformer/vorc_transformer.cpp index 2004ce21d5a574..7ddab79dc1c6b8 100644 --- a/be/src/format/transformer/vorc_transformer.cpp +++ b/be/src/format/transformer/vorc_transformer.cpp @@ -391,11 +391,13 @@ Status VOrcTransformer::collect_file_statistics_after_close(TIcebergColumnStats* continue; } // ORC IDs are depth-first, so top-level fields after a complex field are not i + 1. - uint64_t orc_col_id = orc_root_type.getSubtype(i)->getColumnId(); - if (orc_col_id >= file_stats->getNumberOfColumns()) { + const uint64_t raw_orc_col_id = orc_root_type.getSubtype(i)->getColumnId(); + if (raw_orc_col_id >= file_stats->getNumberOfColumns()) { continue; } + // The uint32_t column-count check above makes narrowing to the ORC API width safe. + const uint32_t orc_col_id = static_cast(raw_orc_col_id); const orc::ColumnStatistics* col_stats = file_stats->getColumnStatistics(orc_col_id); if (col_stats == nullptr) { continue; From d6f2112c00941454deed99c2cf016c1c71b00489 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Mon, 20 Jul 2026 09:53:49 +0800 Subject: [PATCH 10/10] fix(iceberg): preserve ORC upper bound fallback --- .../iceberg/helper/IcebergWriterHelper.java | 19 +++++++----- .../helper/IcebergWriterHelperTest.java | 30 +++++++++++++++++++ 2 files changed, 42 insertions(+), 7 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java index 094cbef038dd03..54a791e7e18133 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java +++ b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelper.java @@ -85,7 +85,7 @@ public static WriteResult convertToWriterResult( long fileSize = commitData.getFileSize(); long recordCount = commitData.getRowCount(); CommonStatistics stat = new CommonStatistics(recordCount, DEFAULT_FILE_COUNT, fileSize); - Metrics metrics = buildDataFileMetrics(commitData, schema, metricsConfig); + Metrics metrics = buildDataFileMetrics(commitData, schema, metricsConfig, fileFormat); Optional partitionData = Optional.empty(); //get and check partitionValues when table is partitionedTable if (spec.isPartitioned()) { @@ -169,7 +169,7 @@ private static PartitionData convertToPartitionData( } private static Metrics buildDataFileMetrics( - TIcebergCommitData commitData, Schema schema, MetricsConfig metricsConfig) { + TIcebergCommitData commitData, Schema schema, MetricsConfig metricsConfig, FileFormat fileFormat) { Map fieldParents = TypeUtil.indexParents(schema.asStruct()); Map columnSizes = new HashMap<>(); Map valueCounts = new HashMap<>(); @@ -201,8 +201,8 @@ private static Metrics buildDataFileMetrics( filterLogicalMetrics(valueCounts, schema, metricsConfig, fieldParents), filterLogicalMetrics(nullValueCounts, schema, metricsConfig, fieldParents), null, - filterBounds(lowerBounds, schema, metricsConfig, fieldParents, true), - filterBounds(upperBounds, schema, metricsConfig, fieldParents, false)); + filterBounds(lowerBounds, schema, metricsConfig, fieldParents, fileFormat, true), + filterBounds(upperBounds, schema, metricsConfig, fieldParents, fileFormat, false)); } private static Map filterDisabledMetrics( @@ -232,7 +232,7 @@ private static Map filterLogicalMetrics( private static Map filterBounds( Map bounds, Schema schema, MetricsConfig metricsConfig, - Map fieldParents, boolean lowerBound) { + Map fieldParents, FileFormat fileFormat, boolean lowerBound) { Map filteredBounds = new HashMap<>(); bounds.forEach((fieldId, value) -> { if (isInRepeatedField(fieldId, schema, fieldParents)) { @@ -248,7 +248,7 @@ private static Map filterBounds( Type type = schema.findType(fieldId); int length = ((MetricsModes.Truncate) mode).length(); // Truncated upper bounds must round up so file pruning cannot exclude matching values. - filteredValue = truncateBound(type, value, length, lowerBound); + filteredValue = truncateBound(type, value, length, fileFormat, lowerBound); } if (filteredValue != null) { filteredBounds.put(fieldId, filteredValue); @@ -269,13 +269,18 @@ private static boolean isInRepeatedField( return false; } - private static ByteBuffer truncateBound(Type type, ByteBuffer value, int length, boolean lowerBound) { + private static ByteBuffer truncateBound( + Type type, ByteBuffer value, int length, FileFormat fileFormat, boolean lowerBound) { switch (type.typeId()) { case STRING: String stringValue = Conversions.fromByteBuffer(type, value).toString(); String truncatedString = lowerBound ? UnicodeUtil.truncateStringMin(stringValue, length) : UnicodeUtil.truncateStringMax(stringValue, length); + // ORC keeps the full maximum when no safe truncated successor exists. + if (!lowerBound && truncatedString == null && fileFormat == FileFormat.ORC) { + return value; + } return truncatedString == null ? null : Conversions.toByteBuffer(type, truncatedString); case BINARY: return lowerBound diff --git a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java index 1769abdcc4c7f6..39fc15ddbe1a7e 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/helper/IcebergWriterHelperTest.java @@ -177,6 +177,36 @@ public void testConvertToWriterResultTruncatesStringAndBinaryBounds() { Assertions.assertEquals(ByteBuffer.wrap(new byte[] {1, 2, 4}), dataFile.upperBounds().get(2)); } + @Test + public void testConvertToWriterResultPreservesOrcUpperBoundWithoutTruncatedSuccessor() { + Schema boundsSchema = new Schema( + Types.NestedField.optional(1, "text", Types.StringType.get())); + Table table = Mockito.mock(Table.class); + Mockito.when(table.schema()).thenReturn(boundsSchema); + Mockito.when(table.spec()).thenReturn(unpartitionedSpec); + Mockito.when(table.sortOrder()).thenReturn(SortOrder.unsorted()); + Mockito.when(table.properties()).thenReturn(Map.of( + TableProperties.DEFAULT_FILE_FORMAT, "orc", + TableProperties.DEFAULT_WRITE_METRICS_MODE, "truncate(1)")); + + String maxWithoutSuccessor = new String(Character.toChars(Character.MAX_CODE_POINT)) + "tail"; + TIcebergColumnStats columnStats = new TIcebergColumnStats(); + columnStats.setUpperBounds(Map.of( + 1, Conversions.toByteBuffer(Types.StringType.get(), maxWithoutSuccessor))); + + TIcebergCommitData commitData = new TIcebergCommitData(); + commitData.setFilePath("/path/to/data.orc"); + commitData.setRowCount(1); + commitData.setFileSize(128); + commitData.setColumnStats(columnStats); + + DataFile dataFile = IcebergWriterHelper.convertToWriterResult(table, List.of(commitData)).dataFiles()[0]; + + Assertions.assertNotNull(dataFile.upperBounds()); + Assertions.assertEquals(maxWithoutSuccessor, Conversions.fromByteBuffer( + Types.StringType.get(), dataFile.upperBounds().get(1)).toString()); + } + @Test public void testConvertToWriterResultBuildsMetricsPolicyOncePerBatch() { Table table = Mockito.mock(Table.class);