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]

Reply via email to