改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

冰山之巅:高效合并小文件的探索之旅

ckckck2025年1月10日20 浏览

1. Spark写入Iceberg,小文件合并

Spark DataFrame 合并

java 复制代码
SparkActions  
  .get(spark)  
  .rewriteDataFiles(icebergTable)  
  .execute()

完整代码

java 复制代码
import org.apache.hadoop.conf.Configuration;  
import org.apache.iceberg.Snapshot;  
import org.apache.iceberg.Table;  
import org.apache.iceberg.actions.Actions;  
import org.apache.iceberg.catalog.Namespace;  
import org.apache.iceberg.catalog.TableIdentifier;  
import org.apache.iceberg.hive.HiveCatalog;  
import org.apache.spark.sql.SparkSession;  
import org.apache.spark.sql.internal.SQLConf;  
import java.util.HashMap;  
import java.util.Map;  
import static org.apache.hadoop.hive.conf.HiveConf.ConfVars.METASTOREURIS;  

public class SparkCompaction {  
    public static void main(String[] args) {  
        val sparkSession: SparkSession = SparkSession.builder()  
                .master("local[*]")  
                .config(SQLConf.PARTITION_OVERWRITE_MODE().key(), "dynamic")  
                .config("spark.hadoop." + METASTOREURIS.varname, "localhost:9083")  
                .config("spark.sql.warehouse.dir", "warehouse")  
                .config("spark.executor.heartbeatInterval", "100000")  
                .config("spark.network.timeoutInterval", "100000")  
                .enableHiveSupport()  
                .getOrCreate();  

        TableIdentifier identifier = TableIdentifier.of(Namespace.of("db"), "table");  
        Map<String, String> config = new HashMap<>();  
        config.put("type", "iceberg");  
        config.put("catalog-type", "hive");  
        config.put("property-version", "1");  
        config.put("warehouse", "warehouse");  
        config.put("uri", "thrift://local:9083");  
        config.put("io-impl", "org.apache.iceberg.aliyun.oss.OSSFileIO");  
        config.put("oss.endpoint", "https://xxx.aliyuncs.com");  
        config.put("oss.access.key.id", "key");  
        config.put("oss.access.key.secret", "secret");  
        HiveCatalog hiveCatalog = new HiveCatalog(new Configuration());  
        hiveCatalog.initialize("iceberg_hive_catalog", config);  
        Table table = hiveCatalog.loadTable(identifier);  

        /**  
         * 合并datafiles小文件,核心代码  
         */  
        SparkActions  
            .get()  
            .rewriteDataFiles(table)  
            .option(BinPackStrategy.MIN_FILE_SIZE_BYTES, "0")  
            .option(RewriteDataFiles.TARGET_FILE_SIZE_BYTES, Long.toString(Long.MAX_VALUE - 1))  
            .option(BinPackStrategy.MAX_FILE_SIZE_BYTES, Long.toString(Long.MAX_VALUE))  
            .option("target-file-size-bytes", (128 * 1024 * 1024).toString) // 128 MB  
            .option("rewrite-job-order", "files-desc")  
            .execute();  

        /**  
         * 合并Manifest小文件,核心代码  
         */  
        SparkActions  
          .get()  
          .rewriteManifests(table)  
          .rewriteIf(manifest -> true)  
          .rewriteIf(file -> file.length() < 10 * 1024 * 1024)  // 10 MB  
          .execute();  

        /**  
         * 写入DataFrame  
         */  
        sortedDf  
            .write()  
            .format("iceberg")  
            .option(SparkWriteOptions.REWRITTEN_FILE_SCAN_TASK_SET_ID, groupID)  
            .option(SparkWriteOptions.TARGET_FILE_SIZE_BYTES, writeMaxFileSize())  
            .option(SparkWriteOptions.USE_TABLE_DISTRIBUTION_AND_ORDERING, "false")  
            .mode("append")  
            .save(groupID);  

        Snapshot snapshot = table.currentSnapshot();  
        if (snapshot != null) {  
            table.expireSnapshots().expireOlderThan(snapshot.timestampMillis()).commit();  
        }  
    }  
}

SparkSQL 存储过程合并小文件

sql 复制代码
-- Rewrite Data Files CALL Procedure in SparkSQL  
CALL catalog.system.rewrite_data_files(  
    table => 'myTable',  
    strategy => 'binpack',  
    options => map(  
        'rewrite-job-order', 'files-desc'  
    )  
);

CALL iceberg_catalog.system.rewrite_data_files(  
    table => 'data_lake_ods.order_info1',  
    options => map(  
        'max-concurrent-file-group-rewrites', '15',  
        'max-file-group-size-bytes', '1073741824',  
        'target-file-size-bytes', '67108864',  
        'rewrite-all', 'true'  
    )  
);

-- 合并方法2大宽表直接合并:  
CALL spark_catalog.system.rewrite_data_files(  
    table => 'data_lake_ods.order_info1',  
    options => map(  
        'max-concurrent-file-group-rewrites', 100,  
        'target-file-size-bytes', '536870912',  
        'partial-progress.enabled', 'true',  
        'rewrite-job-order', 'bytes-asc',  
        'partial-progress.max-commits', '10000',  
        'max-file-group-size-bytes', '10737418240',  
        'rewrite-all', 'true'  
    )  
);

合并小文件策略解析

1)Binpack

sql 复制代码
CALL catalog.system.rewrite_data_files(  
    table => 'streamingtable',  
    strategy => 'binpack',  
    where => 'created_at between "2023-01-26 09:00:00" and "2023-01-26 09:59:59"',  
    options => map(  
        'rewrite-job-order', 'bytes-asc',  
        'target-file-size-bytes', '1073741824',  
        'max-file-group-size-bytes', '10737418240',  
        'partial-progress-enabled', 'true'  
    )  
);

2)Sort

sql 复制代码
CALL catalog.system.rewrite_data_files(  
    table => 'nfl_teams',  
    strategy => 'sort',  
    sort_order => 'team ASC NULLS LAST, name ASC NULLS FIRST'  
);

3)Z-Order

sql 复制代码
CALL catalog.system.rewrite_data_files(  
    table => 'people',  
    strategy => 'sort',  
    sort_order => 'zorder(age,height)'  
);

通过Z-Order策略,满足多字段联合查询提升查询性能的场景。Z-Order本质上是 牺牲了主排序键的查询性能换取次排序键查询性能的提升

2. Flink写入Iceberg,小文件合并

java 复制代码
import org.apache.flink.api.java.utils.ParameterTool;  
import org.apache.hadoop.conf.Configuration;  
import org.apache.iceberg.Snapshot;  
import org.apache.iceberg.Table;  
import org.apache.iceberg.catalog.Catalog;  
import org.apache.iceberg.catalog.Namespace;  
import org.apache.iceberg.catalog.TableIdentifier;  
import org.apache.iceberg.flink.CatalogLoader;  
import org.apache.iceberg.flink.actions.Actions;  
import java.util.HashMap;  
import java.util.Map;  

public class FlinkCompaction {  
    public static void main(String[] args) throws Exception {  
        ParameterTool tool = ParameterTool.fromArgs(args);  
        Map<String, String> properties = new HashMap<>();  
        properties.put("type", "iceberg");  
        properties.put("catalog-type", "hive");  
        properties.put("property-version", "1");  
        properties.put("warehouse", tool.get("warehouse"));  
        properties.put("uri", tool.get("uri"));  
        if (tool.has("oss.endpoint")) {  
            properties.put("io-impl", "org.apache.iceberg.aliyun.oss.OSSFileIO");  
            properties.put("oss.endpoint", tool.get("oss.endpoint"));  
            properties.put("oss.access.key.id", tool.get("oss.access.key.id"));  
            properties.put("oss.access.key.secret", tool.get("oss.access.key.secret"));  
        }  

        CatalogLoader loader = CatalogLoader.hive(tool.get("catalog"), new Configuration(), properties);  
        Catalog catalog = loader.loadCatalog();  
        TableIdentifier identifier = TableIdentifier.of(Namespace.of(tool.get("db")), tool.get("table"));  
        Table table = catalog.loadTable(identifier);  

        /**  
         * 合并小文件,核心代码  
         */  
        Actions.forTable(table)  
                .rewriteDataFiles()  
                .maxParallelism(5)  
                .targetSizeInBytes(128 * 1024 * 1024)  
                .execute();  

        Snapshot snapshot = table.currentSnapshot();  
        if (snapshot != null) {  
            table.expireSnapshots().expireOlderThan(snapshot.timestampMillis()).commit();  
        }  
    }  
}

3. 使用Amoro 来解决小文件问题

Amoro 一个核心功能是 Self-Optimizing。流计算场景下,数据湖表的治理成为核心需求。随着流式写入,它必然带来小文件问题,进而带来查询的性能问题,甚至有可能因为小文件或者 delete 文件太多,造成表不可用。Amoro 数据湖表提供了开箱即用的治理的能力。Amoro 会自动发现注册在 Catalog 中的数据湖表,自动地持续地去监测数据服务表的小文件以及 data 文件的数量,然后对整个表中的文件数据量进行评估,自动根据优先级优化任务调度,达到开箱即用的对数据湖表的文件治理。

Self-Optimizing 核心特点

  • 自动化、异步、透明 — 后台持续检测文件变化,异步分布式执行优化任务,对用户透明、无感知。
  • 资源隔离与共享 — 允许在表级别隔离和共享资源,以及设置资源配额。
  • 灵活且可扩展的部署 — 优化器支持多种部署方式和便捷的扩展。

Self-Optimizing 类型

  1. Minor Optimize:从 fragment file 到 segment file 的转换,把碎片化的小文件转换成非碎片化的文件,合并成大于 16MB 的文件,同时最重要的是做出了这种 eq-delete 操作,它会使 segment file 上面只关联 pos-delete。大幅提升了 Iceberg 表的查询效率。
  2. Major Optimize:做 segment file 内部的转化,部分消除 pos-delete 文件,进一步去碎片化。
  3. Full Optimize:根据用户的配置,比如一天执行一次,会直接把表上的文件合并成期望的 target-size 设置的大小,最终会消除所有的 delete 文件,全部转成 data file。

Self-Optimizing 资源管理

Amoro 通过 Optimizer Group 实现资源隔离和共享。每个 Iceberg 表均属于某个 group,Group 间计算资源相互隔离,互不影响。Group 内部,通过 Quota 分配每个表可以占有的计算资源比例,达到计算资源在不同表之间的隔离或者共享。

Self-Optimizing 部署

Amoro 支持多种部署方式,包括 Local Container 和 Flink Container。Local Container 是本地的基于 JVM 进程的一个 optimize 的进程。Flink Container 是以 Flink 任务的方式去提交,可以当作一个 Flink 集群来使用。Amoro 还支持 Yarn 和 K8s 两种最通用的计算集群,并提供了 External Container,可以通过外部注册的方式,向 AMS 直接注册用户自定义的计算集群类型。


大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...

  • 共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
    • 一、 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
    • 二、免费答疑解惑、交流技术
    • 三、面试指导、模拟面试
    • 四、各类pdf文档下载、星球代码下载
    • 五、提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
    • 另外说明加入星球后支持三天无理由退款,不满意无条件随时退。
    • 需要资料请加微信:D1435221412

end