This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 3ad8169b87f branch-4.1:[bug](iceberg) use split file format for
iceberg scan (#65760) (#65807)
3ad8169b87f is described below
commit 3ad8169b87fecce58600d92e41d3c2d07d63da7c
Author: zhangstar333 <[email protected]>
AuthorDate: Thu Jul 23 14:26:06 2026 +0800
branch-4.1:[bug](iceberg) use split file format for iceberg scan (#65760)
(#65807)
cherry-pick from https://github.com/apache/doris/pull/65760
---
.../create_preinstalled_scripts/iceberg/run31.sql | 33 +++++++++++++
.../apache/doris/datasource/FileQueryScanNode.java | 5 +-
.../doris/datasource/iceberg/IcebergUtils.java | 26 ++--------
.../datasource/iceberg/source/IcebergScanNode.java | 25 +++++-----
.../datasource/iceberg/source/IcebergSplit.java | 3 ++
.../doris/datasource/iceberg/IcebergUtilsTest.java | 22 +++++++++
.../iceberg/source/IcebergScanNodeTest.java | 23 +++++++++
.../iceberg/test_iceberg_mixed_file_formats.out | 12 +++++
.../iceberg/test_iceberg_mixed_file_formats.groovy | 56 ++++++++++++++++++++++
9 files changed, 169 insertions(+), 36 deletions(-)
diff --git
a/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run31.sql
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run31.sql
new file mode 100644
index 00000000000..b1cf79240bb
--- /dev/null
+++
b/docker/thirdparties/docker-compose/iceberg/scripts/create_preinstalled_scripts/iceberg/run31.sql
@@ -0,0 +1,33 @@
+-- Bootstrap an Iceberg table whose active snapshot contains both Parquet and
ORC data files.
+-- This script is sourced on every Iceberg container start, so keep it
repeatable.
+
+CREATE DATABASE IF NOT EXISTS demo.test_db;
+USE demo.test_db;
+
+DROP TABLE IF EXISTS mixed_file_format;
+
+CREATE TABLE mixed_file_format (
+ id INT,
+ source STRING
+)
+USING iceberg
+TBLPROPERTIES (
+ 'format-version' = '2',
+ 'write.format.default' = 'parquet'
+);
+
+-- The first snapshot's data files are Parquet.
+INSERT INTO mixed_file_format VALUES
+ (1, 'parquet'),
+ (2, 'parquet'),
+ (3, 'parquet');
+
+-- Change only the format for subsequent writes. The Parquet files above remain
+-- referenced by the current snapshot, while this append produces ORC files.
+ALTER TABLE mixed_file_format
+ SET TBLPROPERTIES ('write.format.default' = 'orc');
+
+INSERT INTO mixed_file_format VALUES
+ (4, 'orc'),
+ (5, 'orc'),
+ (6, 'orc');
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
index 6f201e5ae66..3bd819bea62 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/FileQueryScanNode.java
@@ -459,8 +459,9 @@ public abstract class FileQueryScanNode extends
FileScanNode {
pathPartitionKeys, partitionValues.getIsNull());
TFileCompressType fileCompressType = getFileCompressType(fileSplit);
rangeDesc.setCompressType(fileCompressType);
- // set file format type, and the type might fall back to native format
in setScanParams
- rangeDesc.setFormatType(getFileFormatType());
+ // Seed connector-specific setup with the scan-level default. A
connector may then
+ // override it with the actual format carried by an individual split.
+ rangeDesc.setFormatType(params.getFormatType());
setScanParams(rangeDesc, fileSplit);
rangeDesc.setFileCacheAdmission(admissionResult);
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 638b0474771..b8b8d596749 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
@@ -1251,7 +1251,7 @@ public class IcebergUtils {
public static FileFormat getFileFormat(Table icebergTable) {
Map<String, String> properties = icebergTable.properties();
- String fileFormatName = resolveFileFormatName(icebergTable,
properties);
+ String fileFormatName = resolveFileFormatName(properties);
FileFormat fileFormat;
if (fileFormatName.toLowerCase().contains(ORC_NAME)) {
fileFormat = FileFormat.ORC;
@@ -1263,7 +1263,7 @@ public class IcebergUtils {
return fileFormat;
}
- private static String resolveFileFormatName(Table icebergTable,
Map<String, String> properties) {
+ private static String resolveFileFormatName(Map<String, String>
properties) {
// 1. Check "write-format" (nickname in Flink and Spark)
if (properties.containsKey(WRITE_FORMAT)) {
return properties.get(WRITE_FORMAT);
@@ -1272,27 +1272,7 @@ public class IcebergUtils {
if (properties.containsKey(TableProperties.DEFAULT_FILE_FORMAT)) {
return properties.get(TableProperties.DEFAULT_FILE_FORMAT);
}
- // 3. Last resort: infer from the actual data files in the current
snapshot.
- // This handles migrated tables where none of the above properties
are set.
- return inferFileFormatFromDataFiles(icebergTable);
- }
-
- private static String inferFileFormatFromDataFiles(Table icebergTable) {
- if (icebergTable.currentSnapshot() == null) {
- LOG.info("Iceberg table {} has no snapshot, defaulting to {}",
icebergTable.name(), PARQUET_NAME);
- return PARQUET_NAME;
- }
- try (CloseableIterable<FileScanTask> files =
icebergTable.newScan().planFiles()) {
- java.util.Iterator<FileScanTask> it = files.iterator();
- if (it.hasNext()) {
- String format = it.next().file().format().name().toLowerCase();
- LOG.info("Iceberg table {} inferred file format {} from data
files", icebergTable.name(), format);
- return format;
- }
- } catch (Exception e) {
- LOG.warn("Failed to infer file format from data files for table
{}, defaulting to {}",
- icebergTable.name(), PARQUET_NAME, e);
- }
+ // Iceberg defaults the write format to Parquet when the table does
not declare one.
return PARQUET_NAME;
}
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
index fac34be4bc0..4e576a8bd67 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergScanNode.java
@@ -23,7 +23,6 @@ import org.apache.doris.analysis.TableSnapshot;
import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.catalog.Env;
import org.apache.doris.catalog.TableIf;
-import org.apache.doris.common.DdlException;
import org.apache.doris.common.UserException;
import org.apache.doris.common.profile.SummaryProfile;
import org.apache.doris.common.security.authentication.ExecutionAuthenticator;
@@ -287,6 +286,8 @@ public class IcebergScanNode extends FileQueryScanNode {
rangeDesc.setTableFormatParams(tableFormatFileDesc);
return;
}
+ // update for every split file format
+
rangeDesc.setFormatType(toTFileFormatType(icebergSplit.getSplitFileFormat()));
if (tableLevelPushDownCount) {
tableFormatFileDesc.setTableLevelRowCount(icebergSplit.getTableLevelRowCount());
} else {
@@ -411,6 +412,15 @@ public class IcebergScanNode extends FileQueryScanNode {
}
}
+ private TFileFormatType toTFileFormatType(FileFormat fileFormat) {
+ if (fileFormat == FileFormat.PARQUET) {
+ return TFileFormatType.FORMAT_PARQUET;
+ } else if (fileFormat == FileFormat.ORC) {
+ return TFileFormatType.FORMAT_ORC;
+ }
+ throw new UnsupportedOperationException("Unsupported Iceberg data file
format: " + fileFormat);
+ }
+
private String getDeleteFileContentType(int content) {
// Iceberg file type: 0: data, 1: position delete, 2: equality delete,
3: deletion vector
switch (content) {
@@ -834,6 +844,7 @@ public class IcebergScanNode extends FileQueryScanNode {
storagePropertiesMap,
new ArrayList<>(),
originalPath);
+ split.setSplitFileFormat(dataFile.format());
if (formatVersion >= 3) {
// -1 means that this table was just upgraded from v2 to v3.
// _row_id and _last_updated_sequence_number column is NULL.
@@ -1045,16 +1056,8 @@ public class IcebergScanNode extends FileQueryScanNode {
if (isSystemTable) {
return TFileFormatType.FORMAT_JNI;
}
- TFileFormatType type;
- String icebergFormat = source.getFileFormat();
- if (icebergFormat.equalsIgnoreCase("parquet")) {
- type = TFileFormatType.FORMAT_PARQUET;
- } else if (icebergFormat.equalsIgnoreCase("orc")) {
- type = TFileFormatType.FORMAT_ORC;
- } else {
- throw new DdlException(String.format("Unsupported format name: %s
for iceberg table.", icebergFormat));
- }
- return type;
+ // for table level file format
+ return toTFileFormatType(IcebergUtils.getFileFormat(icebergTable));
}
@Override
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergSplit.java
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergSplit.java
index 3af484abd6d..2987b12b1d5 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergSplit.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/datasource/iceberg/source/IcebergSplit.java
@@ -23,6 +23,7 @@ import
org.apache.doris.datasource.property.storage.StorageProperties;
import lombok.Data;
import org.apache.iceberg.DeleteFile;
+import org.apache.iceberg.FileFormat;
import java.util.ArrayList;
import java.util.Collections;
@@ -52,6 +53,8 @@ public class IcebergSplit extends FileSplit {
private Long firstRowId = null;
private Long lastUpdatedSequenceNumber = null;
private String serializedSplit;
+ // maybe mixed file format type in one table. so need record it for every
split
+ private FileFormat splitFileFormat;
// File path will be changed if the file is modified, so there's no need
to get modification time.
public IcebergSplit(LocationPath file, long start, long length, long
fileLength, String[] hosts,
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
index 5c669623c41..47011cafc94 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/IcebergUtilsTest.java
@@ -69,6 +69,28 @@ import java.util.Optional;
import java.util.UUID;
public class IcebergUtilsTest {
+ @Test
+ public void testGetFileFormatUsesPropertiesWithoutPlanningDataFiles() {
+ Table table = Mockito.mock(Table.class);
+ Mockito.when(table.properties()).thenReturn(Collections.emptyMap());
+
Mockito.when(table.currentSnapshot()).thenReturn(Mockito.mock(Snapshot.class));
+
+ Assert.assertEquals(org.apache.iceberg.FileFormat.PARQUET,
IcebergUtils.getFileFormat(table));
+ // Do not call newScan planFiles()
+ Mockito.verify(table, Mockito.never()).newScan();
+ }
+
+ @Test
+ public void testGetFileFormatUsesConfiguredTableFormat() {
+ Table table = Mockito.mock(Table.class);
+ Mockito.when(table.properties()).thenReturn(
+ ImmutableMap.of(TableProperties.DEFAULT_FILE_FORMAT, "orc"));
+
+ Assert.assertEquals(org.apache.iceberg.FileFormat.ORC,
IcebergUtils.getFileFormat(table));
+ // Do not call newScan planFiles()
+ Mockito.verify(table, Mockito.never()).newScan();
+ }
+
@Test
public void testParseTableName() {
try {
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
index 46abdf08a58..87e2fb62e52 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/datasource/iceberg/source/IcebergScanNodeTest.java
@@ -24,10 +24,12 @@ import org.apache.doris.datasource.TableFormatType;
import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.thrift.TFileFormatType;
import org.apache.doris.thrift.TFileRangeDesc;
import org.apache.doris.thrift.TIcebergDeleteFileDesc;
import org.apache.iceberg.DataFile;
+import org.apache.iceberg.FileFormat;
import org.apache.iceberg.FileScanTask;
import org.apache.iceberg.Schema;
import org.apache.iceberg.Snapshot;
@@ -137,6 +139,7 @@ public class IcebergScanNodeTest {
IcebergSplit split = new IcebergSplit(LocationPath.of(dataPath), 0,
128, 128, new String[0],
3, Collections.emptyMap(), new ArrayList<>(), dataPath);
split.setTableFormatType(TableFormatType.ICEBERG);
+ split.setSplitFileFormat(FileFormat.PARQUET);
split.setFirstRowId(10L);
split.setLastUpdatedSequenceNumber(20L);
split.setDeleteFileFilters(Collections.emptyList(),
Collections.singletonList(
@@ -158,6 +161,25 @@ public class IcebergScanNodeTest {
Assert.assertEquals((long) Integer.MAX_VALUE + 7L,
deleteFileDesc.getContentSizeInBytes());
}
+ @Test
+ public void testSetIcebergParamsUsesSplitFileFormat() throws Exception {
+ TestIcebergScanNode node = new TestIcebergScanNode(new
SessionVariable());
+ String dataPath = "file:///tmp/data-file.orc";
+ IcebergSplit split = new IcebergSplit(LocationPath.of(dataPath), 0,
128, 128, new String[0],
+ 2, Collections.emptyMap(), new ArrayList<>(), dataPath);
+ split.setTableFormatType(TableFormatType.ICEBERG);
+ split.setSplitFileFormat(FileFormat.ORC);
+
+ Method method =
IcebergScanNode.class.getDeclaredMethod("setIcebergParams",
+ TFileRangeDesc.class, IcebergSplit.class);
+ method.setAccessible(true);
+
+ TFileRangeDesc rangeDesc = new TFileRangeDesc();
+ method.invoke(node, rangeDesc, split);
+
+ Assert.assertEquals(TFileFormatType.FORMAT_ORC,
rangeDesc.getFormatType());
+ }
+
@Test
public void testSetIcebergParamsPropagatesPositionDeleteFileFormat()
throws Exception {
SessionVariable sv = new SessionVariable();
@@ -172,6 +194,7 @@ public class IcebergScanNodeTest {
IcebergSplit split = new IcebergSplit(LocationPath.of(dataPath), 0,
128, 128, new String[0],
2, Collections.emptyMap(), new ArrayList<>(), dataPath);
split.setTableFormatType(TableFormatType.ICEBERG);
+ split.setSplitFileFormat(FileFormat.PARQUET);
split.setDeleteFileFilters(Collections.emptyList(),
Collections.singletonList(
new IcebergDeleteFileFilter.PositionDelete(deletePath, -1L,
-1L, 256L,
org.apache.iceberg.FileFormat.ORC)));
diff --git
a/regression-test/data/external_table_p0/iceberg/test_iceberg_mixed_file_formats.out
b/regression-test/data/external_table_p0/iceberg/test_iceberg_mixed_file_formats.out
new file mode 100644
index 00000000000..42771b3c4e3
--- /dev/null
+++
b/regression-test/data/external_table_p0/iceberg/test_iceberg_mixed_file_formats.out
@@ -0,0 +1,12 @@
+-- This file is automatically generated. You should know what you did if you
want to edit this
+-- !mixed_file_formats_files --
+orc 3
+parquet 3
+
+-- !mixed_file_formats_data --
+1 parquet
+2 parquet
+3 parquet
+4 orc
+5 orc
+6 orc
diff --git
a/regression-test/suites/external_table_p0/iceberg/test_iceberg_mixed_file_formats.groovy
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_mixed_file_formats.groovy
new file mode 100644
index 00000000000..0dfa10f0fb1
--- /dev/null
+++
b/regression-test/suites/external_table_p0/iceberg/test_iceberg_mixed_file_formats.groovy
@@ -0,0 +1,56 @@
+// 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.
+
+suite("test_iceberg_mixed_file_formats",
"p0,external,iceberg,external_docker,external_docker_iceberg") {
+ String enabled = context.config.otherConfigs.get("enableIcebergTest")
+ if (enabled == null || !enabled.equalsIgnoreCase("true")) {
+ logger.info("disable iceberg test.")
+ return
+ }
+
+ String catalogName = "test_iceberg_mixed_file_formats"
+ String restPort = context.config.otherConfigs.get("iceberg_rest_uri_port")
+ String minioPort = context.config.otherConfigs.get("iceberg_minio_port")
+ String externalEnvIp = context.config.otherConfigs.get("externalEnvIp")
+
+ sql """drop catalog if exists ${catalogName}"""
+ sql """
+ CREATE CATALOG ${catalogName} PROPERTIES (
+ 'type' = 'iceberg',
+ 'iceberg.catalog.type' = 'rest',
+ 'uri' = 'http://${externalEnvIp}:${restPort}',
+ 's3.access_key' = 'admin',
+ 's3.secret_key' = 'password',
+ 's3.endpoint' = 'http://${externalEnvIp}:${minioPort}',
+ 's3.region' = 'us-east-1'
+ )
+ """
+ sql """switch ${catalogName}"""
+ sql """ set parallel_pipeline_task_num = 1; """
+ order_qt_mixed_file_formats_files """
+ SELECT lower(file_format), sum(record_count)
+ FROM test_db.mixed_file_format\$files
+ GROUP BY lower(file_format)
+ ORDER BY 1
+ """
+
+ order_qt_mixed_file_formats_data """
+ SELECT id, source
+ FROM test_db.mixed_file_format
+ ORDER BY id
+ """
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]