淘天集团基于 Fluss、Paimon 与 StarRocks 构建湖流一体数据链路

📅 发布时间:2026/8/6 6:59:15
淘天集团基于 Fluss、Paimon 与 StarRocks 构建湖流一体数据链路 在实时业务持续增长、数据决策周期不断缩短的背景下传统实时链路与湖仓链路分离的架构正在越来越难以满足秒级分析需求。对淘天集团这样的大规模业务场景而言核心挑战可以概括为两个方面实时数据产生很快但难以进入统一、可治理、可分析的数据体系。湖仓数据治理和分析能力更强但数据新鲜度通常只能达到分钟级甚至更长延迟。如何让秒级实时数据、分钟级湖仓数据和离线历史数据在同一套体系中连续流转成为我们建设湖流一体链路时首先要解决的问题。一、从“实时链路孤岛”到湖流一体在原有架构中业务日志服务器和业务数据库产生的数据主要进入三类链路秒级实时链路主要依赖 TimeTunnel其定位类似 Kafka负责承载实时消息流。分钟级近实时链路主要使用 Paimon 表实时分支再由 Flink 任务进行流批一体计算并逐层写入 DWD、ADS 或 DWS 等湖仓分层。离线批处理链路主要使用 Paimon 表离线分支StarRocks 在数据服务层读取 Paimon 表为分钟级数据和历史数据提供 OLAP 查询能力。图1 当前湖仓架构这套架构已经能够较好支撑分钟级和离线分析但秒级实时链路与湖仓链路之间仍存在明显断点。TimeTunnel 中的数据通常以字符串形式存在对下游并不直接可见业务、BI 和运营团队如果想直接使用秒级数据往往还需要额外导入到离线表中解析和加工。更重要的是在大促等对时效性要求极高的场景中活动开始后的几秒内业务方就希望看到核心指标变化以便快速判断策略效果和业务风险分钟级数据延迟在这类场景下已经不够及时。图2 业务诉求与核心痛点二、三个组件的分工Fluss、Paimon 与 StarRocks围绕这一目标我们引入 Fluss并将其与 Paimon、StarRocks 结合构建湖流一体数据链路。三者分工如下Fluss负责秒级实时存储和消费提供日志表、主键表、列裁剪和多级分区裁剪等能力。Paimon负责湖仓侧统一存储承担分钟级数据沉淀、历史数据管理和流批一体计算结果承载。StarRocks负责统一查询服务将 Fluss 秒级增量与 Paimon 分钟级、历史数据纳入同一条 OLAP 分析路径并复用其在 Paimon 表扫描、谓词下推和复杂查询优化上的能力为业务提供高性能查询入口。在新的架构中秒级实时数据、分钟级湖仓数据和离线历史数据通过统一查询入口衔接起来分层承载ODS、DWD、ADS/DWS 等数据层级均可以使用 Fluss 承载秒级实时数据。自动沉淀开启湖流一体能力后Fluss 会自动启动分层同步服务将实时数据周期性同步到 Paimon 表中形成分钟级湖仓数据离线批处理数据继续写入 Paimon 历史分支。统一查询服务层仍以 StarRocks 作为统一查询入口。需要秒级新鲜度时访问 Fluss 中的最新增量数据访问分钟级和历史数据时查询 Paimon 表。这样业务不需要理解底层链路差异也能在统一查询路径中获得更完整的数据时效性。在这一链路中StarRocks 的价值不只是“读表”而是把 Fluss 的秒级增量数据和 Paimon 的分钟级、历史数据组织成统一的分析视图。业务侧面对的是一条 OLAP 查询路径底层则由 StarRocks 根据数据新鲜度和同步进度选择更合适的数据访问位置从而避免业务在实时表、湖仓表和历史表之间手工拼接口径。图3 湖流一体数据架构三、 Fluss 承接实时链路高吞吐、更新语义与低成本消费在实践中Fluss 主要通过日志表和主键表承接不同实时数据场景。两类表分别面向不同语义日志表面向 append-only 场景按照写入顺序存储数据使用方式类似 Kafka 或 TimeTunnel。它不支持更新和删除但具备较高写入吞吐实践中峰值写入能力达到每秒 3500 万条。主键表面向需要最新状态的数据支持 INSERT、UPDATE 和 DELETE 操作适用于交易、订单、用户状态等更新语义场景实践中峰值写入能力达到每秒 500 万条。主键表的 LastRow 引擎还支持部分列更新在多路实时流共同写入一张宽表时新写入的非空字段会更新结果表未写入或为空的字段不会影响已有值。这使得用户流、订单流等多路数据可以以同一主键汇入一张实时宽表显著降低多流合并和宽表构建的复杂度。图4 主键表部分列更新除了表模型本身Fluss 在消费成本优化上的收益也非常关键主要体现在列裁剪和多级分区裁剪两个方面列裁剪解决“读取哪些字段”的问题。下游任务只消费 id、status、省份、平台等少量字段时如扩展信息、参数等大字段不会被读取也不会参与反序列化部分场景消费带宽可降低 90% 以上。图5 列裁剪多级分区裁剪解决“读取哪些行”的问题。下游任务可以只读取目标分区中的数据而不是先消费全量数据再过滤。图6 多级分区裁剪多级分区设计面向分区值集中的场景建议使用等值进行过滤面向分区值分散的场景建议写入时通过Hash算法将其映射到固定的桶中可以有效避免数据倾斜。图7 多级分区设计四、湖流同步与 StarRocks 查询兼顾实时性和分析性能湖流一体链路的核心在于实时数据的自动沉淀。Fluss 与 Paimon 分别承担不同的数据存储职责Fluss 表采用流式 Arrow 格式存储适合低延迟读写和短期实时数据保留。Paimon 表采用 Parquet 格式存储压缩率更高适合分钟级和历史数据分析。自动同步开启湖流一体后Fluss 会自动启动分层服务将秒级数据定期同步到 Paimon 实时分支离线批处理数据继续写入 Paimon 历史分支。图8 湖流一体链路搭建在配置上湖流一体能力通过表参数开启核心参数包括table.datalake.enabled true开启湖流一体能力Fluss 会自动创建字段结构和表路径一致的 Paimon 表。table.datalake.freshness控制 Fluss 写入 Paimon 的频率默认值为 3 分钟可根据业务实时性要求调整。paimon.前缀参数用于指定 Paimon 表属性例如通过paimon.file.format配置 Paimon 表文件格式。图9 湖流一体表参数设置StarRocks 在湖流一体查询中承担统一分析入口。Fluss 每次同步到 Paimon 时都会产生 checkpointStarRocks 可以据此将一次查询拆分为两段数据访问checkpoint 之前的数据访问 Paimon复用 StarRocks 读取 Paimon 表的既有优化能力覆盖绝大多数历史数据和分钟级数据。checkpoint 之后的增量数据访问 Fluss只查询最近几分钟的秒级增量因此对整体查询延迟影响有限。面向业务的结果底层路径可以分段执行但对业务侧呈现为一条统一的 OLAP 查询路径。这种分段读取方式也是 StarRocks 在湖流一体链路中的关键优势大部分历史数据继续走 Paimon 查询路径可以复用 StarRocks 对湖仓表扫描、列裁剪、谓词下推和复杂 OLAP 计算的优化只有 checkpoint 之后的少量最新数据访问 Fluss从而把秒级新鲜度的成本控制在较小范围内。换句话说StarRocks 将“读湖仓的高吞吐”和“读实时流的低延迟”组合在同一条查询链路中使实时分析不再依赖额外的数据搬运或人工拼接。图10 基于湖流一体链路的 OLAP 查询五、阶段成果与后续演进经过阶段性建设湖流一体链路已经从架构验证进入实际应用阶段。其价值不仅体现在引入新组件更体现在将秒级实时数据纳入统一数据体系实时数据可见湖仓数据可复用分析查询路径更统一实时链路成本也得到显著降低。秒级数据使用门槛下降具备基础 SQL 能力的业务用户可以直接通过 StarRocks 查询 Fluss 表完成简单实时分析。OLAP 查询稳定可用StarRocks 读取 Fluss 表并结合 Paimon 历史数据后可以支持秒级响应的 OLAP 查询。开发运维效率提升秒级场景开发和运维效率提升 50% 以上开发验证周期从 5 天缩短到 2 天。链路成本显著降低借助 Fluss 的列裁剪和多级分区裁剪能力消费带宽和反序列化总成本减少 80% 以上。后续我们会重点沿三个方向继续推进集团内推广 Fluss扩大列裁剪和多级分区裁剪能力的应用范围持续降低实时数据消费成本。建设实时物化视图当 Fluss 中数据发生变化时StarRocks 端能够自动、增量地更新物化视图避免全表扫描。探索 AI 实时规则引擎将 Fluss 秒级数据流作为 AI Agent 和规则系统的输入在关键业务指标异常时实时识别并推送给业务方。总体而言淘天集团的湖流一体实践可以概括为一条主线以 Fluss 补齐秒级实时数据的 Schema 化和低成本消费能力以 Paimon 承接分钟级与历史数据沉淀以 StarRocks 提供统一、高性能的 OLAP 查询能力。三者结合后StarRocks 不只是服务层查询入口更是打通 Fluss 秒级增量与 Paimon 湖仓数据的分析引擎使秒级数据、分钟级数据和历史数据可以在统一体系中被管理和分析。这套架构的意义不只是提升单个场景的实时性而是为数据平台提供了一种更连续的数据组织方式。未来随着实时物化视图和 AI 实时规则引擎的建设StarRocks 在湖流一体链路中的角色还会从统一查询入口进一步扩展到实时预计算、自动增量更新和面向业务决策的分析服务底座。六、阿里云 EMR Serverless StarRocks 对 Fluss 的适配与增强阿里云 EMR Serverless StarRocks 对 Fluss 进行了全面的原生适配提供从 Catalog 注册、数据读取、分区裁剪到 Union Read 的完整支持并叠加了商业版独有的性能增强。该方案旨在通过“湖流一体”架构替代传统的 Kafka Flink ETL 数据湖复杂链路实现开箱即用、低成本且高性能的实时数据分析。能力说明商业版增强Fluss Catalog通过CREATE EXTERNAL CATALOG注册 Fluss 数据源SQL 直查预置 Catalog 模板一键配置Native 读取 Paimon 数据对 Fluss 湖侧Paimon 格式数据实现原生 C 读取绕过 JNI 开销Stella 自研算子性能领先开源分区裁剪按分区条件过滤 Fluss 表避免全表扫描与内表一致的裁剪优化Union Read一条 SQL 自动合并 Fluss 实时数据 Paimon 历史数据全托管无需运维 Tiering ServiceNative SDK 接入StarRocks 通过 Fluss Native SDK 直连替代 JNI 调用链路进一步降低读取开销提升吞吐读取性能持续优化缩小与内表查询性能差距目标 1.5x 以内湖流一体场景下查询性能对标内表EMR Serverless StarRocks 依托自研 Stella 引擎在 Fluss 场景下具备三项开源不具备的核心优化Native Reader 性能优化对 Fluss Lake Split 实现原生 C 读取结合向量化批处理与列式裁剪大幅减少反序列化开销显著提升查询性能。DataCache 统一管理Fluss 外表查询与 StarRocks 内表共享统一 DataCache 层热数据自动缓存至本地 SSD重复查询直接命中避免了内外表缓存空间冲突问题。查询优化器增强支持 Fluss 表与内表 Join时的最优策略选择以及跨 CatalogFluss Paimon 内表的全局查询优化。相比传统“Kafka Flink ETL 数据湖 StarRocks”架构新方案在多个维度具有显著优势维度传统方案 (KafkaFlinkLakeSR)新方案 (Fluss EMR Serverless SR)架构复杂度4套系统各自运维协调困难2套系统全托管极简架构数据存储Kafka 数据湖双写存储成本高Fluss 单份数据流湖视图分离存储成本低ETL 链路需开发维护 Flink 作业导入数据湖内置 Tiering Service自动同步零 ETL 代码查询实时性分钟级受限于 ETL 批次间隔秒级Union Read 直查实时层数据一致性需人工对齐 Kafka 与湖数据口径原生一致同一份数据源运维/TCO高运维成本高存储成本低运维AI Agent 智能运维低 TCOServerless 按需付费 单写存储