This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 54e0e9712b6 [bug](iceberg) use split file format for iceberg scan
(#65760)
54e0e9712b6 is described below
commit 54e0e9712b6057d5f1c643103d1612d196b3cd35
Author: zhangstar333 <[email protected]>
AuthorDate: Sat Jul 18 19:34:11 2026 +0800
[bug](iceberg) use split file format for iceberg scan (#65760)
### What problem does this PR solve?
Problem Summary:
A table maybe contain mixed parquet and orc files.
before used the table-level format for every split files, it's should
use from data file.
and when table properties were absent, maybe trigger full planFiles()
call for every splits, this cause mush times.
---
.../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 | 22 +++++++++
.../iceberg/test_iceberg_mixed_file_formats.out | 12 +++++
.../iceberg/test_iceberg_mixed_file_formats.groovy | 56 ++++++++++++++++++++++
9 files changed, 168 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 e510f05bfb6..b1dc260701e 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
@@ -480,8 +480,9 @@ public abstract class FileQueryScanNode extends
FileScanNode {
pathPartitionKeys, partitionValuesFromPath.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 12b5dce1139..b6f1d22046d 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
@@ -1293,7 +1293,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;
@@ -1305,7 +1305,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);
@@ -1314,27 +1314,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 b592883f860..972009fb1e7 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
@@ -25,7 +25,6 @@ import org.apache.doris.analysis.TupleDescriptor;
import org.apache.doris.catalog.Column;
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;
@@ -329,6 +328,8 @@ public class IcebergScanNode extends FileQueryScanNode {
rangeDesc.unsetColumnsFromPathIsNull();
return;
}
+ // update for every split file format
+
rangeDesc.setFormatType(toTFileFormatType(icebergSplit.getSplitFileFormat()));
if (tableLevelPushDownCount) {
tableFormatFileDesc.setTableLevelRowCount(icebergSplit.getTableLevelRowCount());
} else {
@@ -498,6 +499,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) {
@@ -1008,6 +1018,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.
@@ -1430,16 +1441,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 59b4f483f01..eeeff694b8e 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
@@ -24,6 +24,7 @@ import org.apache.doris.thrift.TFileFormatType;
import lombok.Data;
import org.apache.iceberg.DeleteFile;
+import org.apache.iceberg.FileFormat;
import java.util.ArrayList;
import java.util.Collections;
@@ -53,6 +54,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;
private boolean positionDeleteSystemTableSplit = false;
private TFileFormatType positionDeleteFileFormat;
private int positionDeleteContent;
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 ee34b4bee50..349e1e0bb82 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
@@ -78,6 +78,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 testGetIcebergViewUsesSessionCatalogWithDelegatedCredential() {
ConnectContext context = new ConnectContext();
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 d4273f6faa4..c823cb2cc89 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
@@ -30,6 +30,7 @@ import org.apache.doris.planner.PlanNodeId;
import org.apache.doris.planner.ScanContext;
import org.apache.doris.qe.SessionVariable;
import org.apache.doris.system.Backend;
+import org.apache.doris.thrift.TFileFormatType;
import org.apache.doris.thrift.TFileRangeDesc;
import org.apache.doris.thrift.TIcebergDeleteFileDesc;
@@ -250,6 +251,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(
@@ -271,6 +273,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();
@@ -285,6 +306,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..ebcb2affa29
--- /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]