ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

存算分离架构实战:从对象存储到数据湖高效整合

存算分离架构实战:从对象存储到数据湖高效整合 存算分离这四个字最近几年在大数据圈子里出现频率高得吓人。我最早理解它是在公司 Hadoop 集群被业务逼到墙角的时候——存储快满了CPU 和内存却很闲想扩容只能加节点一加就是几十台计算和存储必须一起买。后来我们把存储整体搬上对象存储计算集群瘦身到原来的三分之一跑批性能和稳定性反而上去了。更重要的是不同团队终于开始共享同一份数据而不是各自维护一套副本。这篇文章就围绕存算分离这条技术路线展开重点讲讲它为什么能解决数据高效整合的问题以及怎么落地。不管你是大数据架构师、数据开发还是正在做技术选型的负责人只要你的团队面临集群资源浪费、数据重复拷贝、跨部门取数困难这些问题这篇文章都值得仔细看。我会从架构演进原理讲起再把数据湖表格式、多引擎协同、元数据体系这些关键点拆开最后给出一个可以在自己环境里复刻的完整实操方案以及我实际踩过的一些坑。1. 从存算一体到存算分离我们为什么非拆不可1.1 Hadoop时代存储和计算的“捆绑销售”十年前的 Hadoop 体系本质上是一个存算一体的分布式系统。HDFS 里的 DataNode 既要负责存数据块也要承载 MapReduce、Spark 的计算任务。这样做最大的收益是数据本地性计算被调度到数据所在的节点上减少网络传输处理几百 GB 的日志时体验非常顺滑。但代价也很大。第一扩容必须计算和存储一起扩展。你只是磁盘不够了也得上几十台带 CPU 内存的机器你只是想加算力跑机器学习也得顺带买一堆用不上的存储空间。第二集群内部资源互相干扰。同一个集群里A 团队跑 SQL 大量扫盘B 团队训模型大量占 CPU两者叠加很容易出现互相拖慢的情况。第三利用率普遍偏低。很多生产集群的 CPU 平均使用率常年不到 30%但磁盘和内存老是告急。还有一个很现实的问题数据是“跟着机器走的”换一套集群、换一个供应商迁移成本高得吓人。所以企业在建设大数据平台时往往不敢轻易动底层最终形成越积越多的“烟囱式”系统。1.2 对象存储的崛起把“拆”变成了可能存算分离真正走向成熟离不开对象存储。S3、OSS、COS 这些对象存储本质上是一个无限扩展、按量付费的存储服务它只提供简单的 GET/PUT/LIST 操作天然和计算解耦。你把数据放进去计算集群想开几台就开几台用完可以缩掉存储本身不动。有人可能会问对象存储不是没有随机写能力吗早期 Hadoop 生态确实不兼容它。真正改变局面的是 Hadoop S3A 文件系统客户端的成熟以及后来数据湖表格式的兴起。Spark、Flink、Trino 都可以通过统一的 S3 协议读写对象存储而数据湖表格式在对象存储之上补上了最关键的 ACID 事务能力。没有这一层存算分离就是空中楼阁。另外像 JuiceFS、Alluxio 这类的缓存加速层可以在对象存储和计算节点之间加一层本地缓存把热数据的读取性能拉高一个数量级。这也是存算分离架构里很常见的一个组件。1.3 存算分离如何顺手解决了数据整合难题企业数据整合最大的痛点从来不是“数据不够多”而是数据散落在各个地方业务库在 MySQL用户在 MongoDB日志在 Kafka指标在 ClickHouse。想做一次跨域分析往往要把数据导来导去每个团队各存一份口径还不一致。存算分离之后所有源系统的数据都可以先入湖到同一个对象存储底座比如统一的 ODS 层。计算层可以有多个引擎各司其职Flink 做实时入湖、Spark 做批量清洗加工、Trino 做即席查询。它们读取的是同一份数据、同一个表格式、同一个元数据服务。这在架构上消灭了“数据复制”的源头数据从产生到使用只有一份物理副本这就是“高效整合”最实在的体现。当然数据整合不等于“放一起就算完”还需要元数据、权限、质量、血缘统一治理。存算分离恰恰给这些治理动作提供了一个集中的落点这是传统分布式文件系统很难做到的。2. 实现数据高效整合的技术地基数据湖表格式与元数据2.1 没有ACID对象存储上就没法谈整合对象存储本身是面向“读写文件”设计的不是面向“读写一张表”设计的。如果多个任务同时往同一个目录写文件查询端很可能会读到一半的数据。更糟糕的是对象存储的 LIST 操作在大目录下性能很差文件一多就卡。没有事务能力的话多个引擎根本不敢同时写一张表。于是 Iceberg、Hudi、Delta Lake 这三个数据湖表格式成了存算分离的基石。以 Iceberg 为例它的机制有点像数据库的 MVCC一张表对应一个 metadata 文件里面记录了当前快照指向哪些 manifest 清单每个清单又指向具体的数据文件。写入数据时新数据先生成一组新文件然后通过原子操作把 metadata 指针切换到新快照。读请求要么看到旧快照要么看到新快照永远不会看到一半。这三个格式各有侧重Delta Lake 和 Spark 生态绑定较深Hudi 在 upsert 和近实时更新上更强Iceberg 对多引擎的兼容性和云存储适配做得最均衡。如果团队没有特殊要求我建议新项目优先考虑 Iceberg社区活跃、Spark/Flink/Trino 都有非常成熟的集成。2.2 Spark、Flink、Trino共享同一份数据的工作机制存算分离架构里计算引擎通常不止一个而它们必须读同一张表。以我常用的组合为例Flink 从 Kafka 消费订单数据实时写入 Iceberg 的 ods_orders 表Spark 每天凌晨跑批量任务把 ods 层数据清洗加工成 dwd 层宽表业务部门用 Trino 直接查 Iceberg 表做报表。这三个引擎为什么不会互相把数据写乱核心还是快照隔离机制。每次提交都是一次独立的快照切换读端按当前快照读写端并发通过乐观锁解决冲突。比如两个 Spark 任务同时提交其中一个会失败并自动重试不会出现脏写。而且 Iceberg 支持时间旅行可以随时按某个快照 ID 或时间点往回查这在传统 Hive 表上基本做不到。多引擎共享还有一个好处实时链路和离线链路可以共用一套表。以前实时报表和离线报表各有一套数据两边时常对不上现在 Flink 实时写入、Spark 按时批量合并BI 查询永远读到同一份存储的最新快照数据口径自然一致。2.3 元数据先行Hive Metastore、Catalog与权限层怎么搭存算分离之后元数据服务的重要性甚至超过存储本身。表在哪、分区怎么切、schema 长什么样、谁有权限读这些信息必须要有一个集中的地方管理。最常见的元数据服务是 Hive Metastore也可以选云上的 Glue Data Catalog或者用 Iceberg 的 REST Catalog。我个人的建议是如果团队已有 Hive 生态先用 Hive Metastore 成本最低几乎零改造如果是从零开始做数据湖平台可以考虑 REST Catalog它更轻量和 Iceberg 绑定更自然。权限这块必须分层设计。存储层用对象存储桶策略限制文件访问元数据层用 Ranger 控制表和列的权限计算层再结合 Trino/Spark SQL 的用户体系做行级过滤和列脱敏。关键一点永远不要给普通用户直接暴露对象存储的 AccessKey否则他会绕开所有数据权限直接读源文件整条权限链等于白搞。3. 实操搭一套MinIOIcebergSparkTrino的存算分离数据平台3.1 架构选型与组件版本如果你要在本机或测试环境搭建一套最小可用的存算分离平台我推荐这个组合存储MinIO兼容 S3 协议本地模拟对象存储最方便生产环境换成云上 S3/OSS。表格式Apache Iceberg负责 ACID、快照管理和小文件合并。元数据Hive Metastore 或者 Iceberg REST Catalog二选一。批量计算Spark 3.x负责清洗、加工、合并。实时计算Flink负责实时写入和流批一体。即席查询Trino负责接 BI 和即席分析。整体链路大概是Kafka - Flink - Iceberg on MinIO - Spark(批量加工) - Trino(查询) - BI。所有引擎都通过同一个 Catalog 指向同一个 Warehouse数据只有一份。需要提醒的是生产环境不要用 MinIO 替代云对象存储。云厂商的对象存储在可用性、带宽、生命周期管理上都要成熟得多而且往往有低频存储和冷归档能显著降低存储成本。MinIO 更适合开发测试或者私有化部署且对合规要求极高的场景。3.2 配置与参数计算先算一组网络带宽。假设 ODS 层每天新增数据 100GB同步窗口 2 小时那理想平均带宽是 100×1024 MB / 7200 秒 ≈ 14.2 MB/s。但真实任务在高峰期经常是平均值的好几倍再加上 Flink 写 checkpoint、Iceberg 提交、Trino 读远端数据也会占带宽我会按 70100 MB/s 来规划内网带宽。计算集群和存储集群必须放同一个可用区走内网访问对象存储千万别走公网。计算节点的规格选择也和传统集群不同。以前 Hadoop 节点要大内存大磁盘现在计算节点只需要大 CPU 大内存本地磁盘只放临时文件、缓存和 shuffle 中间结果。一般按 CPU 与内存 1:2 到 1:4 的配比选计算型实例比如 64C256G 这种。每台节点挂一块 几百GB 到 1TB 的 SSD 做缓存就够。Spark 连接 S3 的关键配置我直接给一段可用的 PySpark 初始化代码from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(iceberg-demo) \ .config(spark.sql.extensions, org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions) \ .config(spark.sql.catalog.demo, org.apache.iceberg.spark.SparkCatalog) \ .config(spark.sql.catalog.demo.type, hadoop) \ .config(spark.sql.catalog.demo.warehouse, s3a://data-lake/warehouse) \ .config(spark.sql.catalog.demo.io-impl, org.apache.iceberg.io.ResolvingFileIO) \ .config(spark.hadoop.fs.s3a.endpoint, http://minio:9000) \ .config(spark.hadoop.fs.s3a.access.key, minioadmin) \ .config(spark.hadoop.fs.s3a.secret.key, minioadmin) \ .config(spark.hadoop.fs.s3a.path.style.access, true) \ .getOrCreate()如果是云上 S3就不需要 endpoint 和 path.style.access直接用 access key 和 secret key 就行。路径前缀 s3a:// 是 Hadoop S3A 文件系统的标准写法。3.3 建表、入湖与数据整合落地配置好 Catalog 之后先建表。用 Trino 或者 Spark 都可以我习惯在 Trino 里建CREATE TABLE iceberg.demo.ods_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(18,2), status VARCHAR, dt DATE ) PARTITIONED BY (dt) WITH (format PARQUET);接着把 Kafka 里的订单数据通过 Flink SQL 实时写入这张 Iceberg 表。Flink SQL 里的建表语句类似这样CREATE CATALOG iceberg_demo WITH ( type iceberg, catalog-type hive, uri thrift://hive-metastore:9083, warehouse s3a://data-lake/warehouse ); CREATE TABLE ods_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(18,2), status VARCHAR, dt DATE ) WITH ( connector kafka, topic ods_orders, properties.bootstrap.servers kafka:9092, format json, scan.startup.mode earliest-offset ); INSERT INTO iceberg_demo.demo.ods_orders SELECT order_id, user_id, amount, status, dt FROM ods_orders;Flink 写入和 Spark 读写之所以能同时进行正是靠 Iceberg 的乐观并发控制。这里要注意不同 Flink Iceberg 版本的建表语法细节有差异以官方文档为准。ODS 层只是把原始数据搬进来真正的数据整合在 DWD 层。比如把订单表和用户维度表做宽表加工INSERT INTO iceberg.demo.dwd_order_user SELECT o.order_id, u.user_name, o.amount, o.dt FROM iceberg.demo.ods_orders o LEFT JOIN iceberg.demo.dim_users u ON o.user_id u.user_id WHERE o.dt CURRENT_DATE;这一步的好处是源系统里分散在 MySQL、MongoDB、Kafka 的数据被统一收拢到同一套 Iceberg 表里加工成宽表后下游查询再也不用自己去 join 多套系统。3.4 统一查询与分析让不同团队都用同一份数据数据整合完成后消费端只要连 Trino 就能查所有表。Trino 的 Iceberg Catalog 配置通常是connector.nameiceberg iceberg.catalog.typehive hive.metastore.urithrift://hive-metastore:9083然后在 BI 工具里连接 Trino 的 JDBC 端口就能直接查询 Iceberg 表。以前实现一个综合报表DBA 每天要导 MySQL、ETL 清洗日志、再导入分析库中间至少三次数据拷贝现在只需要在 Trino 里写一条 SQL把订单宽表、用户维表、行为聚合表关联起来就行。实时链路和离线链路也终于能在同一张表上对齐口径这是我在实际业务里感知最强的收益。4. 存算分离落地会踩的坑问题排查与运维经验4.1 常见故障速查表先给一张我在实际运维里整理的速查表基本覆盖了存算分离平台最常见的坑现象可能原因排查与解决写入后其他引擎查不到新数据Catalog 指向不一致或元数据缓存未刷新检查所有引擎的 Catalog 配置是否指向同一份 Warehouse 和元数据服务查询远程存储特别慢数据本地性失效、小文件多、未做谓词下推开启本地缓存合理分区定期合并小文件尽量使用列式过滤高频流式提交后小文件爆炸Flush 间隔太短、写入并发太高调大 checkpoint 间隔和 target-file-size定期调用 rewrite_data_files并发写表报 Commit conflict两个任务同时基于旧快照提交开启 Iceberg 的提交重试机制降低写并发按分区拆分任务Spark 写 S3 任务长时间卡住Hadoop 文件提交器兼容性问题走 Iceberg 的写入路径不要直接依赖 S3A Committer检查线程池和 JMX 指标对象存储 Request 费用暴涨小文件太多导致大量 GET/LIST 请求合并文件、加缓存减少高频提交开启生命周期策略4.2 三个容易被忽略的深层问题第一个是小文件问题它不是一次性任务而是持续过程。流式任务每 5 分钟 flush 一次一天就是 288 个小文件一个月就是八千多个文件对象存储的 LIST 性能会被拖死。我的做法是流式任务把 checkpoint 间隔拉长到 10 分钟以上同时设置 Iceberg 的 write.target-file-size-bytes 为 128MB每天凌晨用 Spark 跑一次重写任务CALL iceberg.system.rewrite_data_files( table demo.db.ods_orders, strategy binpack, options map( min-input-files, 5, target-file-size-bytes, 134217728 ) );第二个是跨引擎 Schema 演进的问题。Iceberg 支持新增字段但老版本引擎可能不认识新字段会导致读取报错。我建议整个平台尽量统一 Iceberg format-version2并升级所有引擎到支持它的版本。上生产之前先用一个测试表做一遍“加字段 - 各引擎读取”的冒烟测试。第三个是成本模型的变化。存算分离之后计算成本从“买机器”变成“按运行时长付费”存储成本则变成“按数据量和请求次数付费”。很多人忽略的是对象存储的 Request 费用大量 Get 和 List 请求可能比存储费还贵。控制费用的核心是缓存和合并给热数据做好缓存、把冷数据及时转低频存储或归档能省下不少成本。最后再分享一个我在实际运维中的体会存算分离不是一个纯技术方案更像是一次数据治理的重启。把数据从分散的集群搬到对象存储只是第一步真正值钱的是后面那套元数据、权限、质量、血缘体系。如果你也被数据重复、集群利用率低、跨部门取数困难这些问题折磨建议先拿一个业务域做最小试点不要一次把全公司数据都搬上去。先让一套链路跑通再逐步扩展你会慢慢感受到从“维护集群”到“维护数据资产”的转变。
返回列表