冰山之巅:高效合并小文件的探索之旅
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 类型
- Minor Optimize:从 fragment file 到 segment file 的转换,把碎片化的小文件转换成非碎片化的文件,合并成大于 16MB 的文件,同时最重要的是做出了这种 eq-delete 操作,它会使 segment file 上面只关联 pos-delete。大幅提升了 Iceberg 表的查询效率。
- Major Optimize:做 segment file 内部的转化,部分消除 pos-delete 文件,进一步去碎片化。
- 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
