作者:朱奥,淘天集团
在实时业务持续增长、数据决策周期不断缩短的背景下,传统实时链路与湖仓链路分离的架构,正在越来越难以满足秒级分析需求。对淘天集团这样的大规模业务场景而言,核心挑战可以概括为两个方面:
- 实时数据产生很快,但难以进入统一、可治理、可分析的数据体系。
- 湖仓数据治理和分析能力更强,但数据新鲜度通常只能达到分钟级甚至更长延迟。
如何让秒级实时数据、分钟级湖仓数据和离线历史数据在同一套体系中连续流转,成为我们建设湖流一体链路时首先要解决的问题。
一、从“实时链路孤岛”到湖流一体
在原有架构中,业务日志服务器和业务数据库产生的数据主要进入三类链路:
- 秒级实时链路:主要依赖 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 时都会产生 checkpoint,StarRocks 可以据此将一次查询拆分为两段数据访问:
- 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 |
通过 |
预置 Catalog 模板,一键配置 |
Native 读取 Paimon 数据 |
对 Fluss 湖侧(Paimon 格式)数据实现原生 C++ 读取,绕过 JNI 开销 |
Stella 自研算子,性能领先开源 |
分区裁剪 |
按分区条件过滤 Fluss 表,避免全表扫描 |
与内表一致的裁剪优化 |
Union Read |
一条 SQL 自动合并 Fluss 实时数据 + Paimon 历史数据 |
全托管,无需运维 Tiering Service |
Native 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时的最优策略选择,以及跨 Catalog(Fluss + Paimon + 内表)的全局查询优化。
相比传统“Kafka + Flink ETL + 数据湖 + StarRocks”架构,新方案在多个维度具有显著优势:
维度 |
传统方案 (Kafka+Flink+Lake+SR) |
新方案 (Fluss + EMR Serverless SR) |
架构复杂度 |
4套系统,各自运维,协调困难 |
2套系统,全托管,极简架构 |
数据存储 |
Kafka + 数据湖双写,存储成本高 |
Fluss 单份数据,流湖视图分离,存储成本低 |
ETL 链路 |
需开发维护 Flink 作业导入数据湖 |
内置 Tiering Service,自动同步,零 ETL 代码 |
查询实时性 |
分钟级(受限于 ETL 批次间隔) |
秒级(Union Read 直查实时层) |
数据一致性 |
需人工对齐 Kafka 与湖数据口径 |
原生一致,同一份数据源 |
运维/TCO |
高运维成本,高存储成本 |
低运维(AI Agent 智能运维),低 TCO(Serverless 按需付费 + 单写存储) |