引言

        在大数据实时处理场景中,订单佣金结算、多维度推广效果分析等需求对数据延迟、查询性能和扩展性提出了极高要求。传统 T+1 离线计算模式已无法满足业务实时性诉求,而单一架构又难以平衡 “实时处理” 与 “复杂查询” 的双重需求。本文结合实际项目实践,分享一套基于 Flink+ClickHouse 的千万级数据实时计算架构,重点拆解 DWD+DWS 分层设计逻辑、核心技术选型与性能优化思路,为同类场景提供可复用的解决方案。

一、业务背景与核心挑战

1. 业务核心诉求

需支撑千万级订单数据的实时处理,提供三类核心服务:

  • 实时收益核算:推广人员、商户需实时查看订单佣金与推广效果;
  • 多维度分析:支持按推广渠道、商品品类、用户等级等维度的交叉查询;
  • 佣金结算支撑:保障核心业务链路的高可用与数据准确性,满足财务结算需求。

2. 传统架构的核心痛点

  • 数据延迟严重:T+1 离线计算模式,数据处理延迟达小时级,无法支撑实时决策;
  • 查询性能瓶颈:海量明细数据导致复杂报表查询耗时长达数小时,用户体验极差;
  • 扩展性不足:新增业务维度或计算规则时,需大量改造核心代码,迭代效率低;
  • 数据一致性差:离线计算易出现数据重复或丢失,影响佣金核算准确性。

二、核心架构设计:Kafka→Flink→ClickHouse 全链路方案

围绕 “实时性、高可用、可扩展” 三大目标,设计 “数据接入 - 实时计算 - 存储查询” 全链路架构,核心采用 DWD+DWS 分层模型,结合 Flink 与 ClickHouse 的协同优势,平衡 “低延迟” 与 “高性能” 的双重需求。

1. 整体架构链路

业务数据(订单/用户行为/维度数据)→ Kafka(数据缓冲分发)→ Flink(实时计算)→ ClickHouse(存储查询)→ 业务应用(报表/结算/告警)
  • 数据接入层:Kafka 承接日均 1 亿 + 用户行为事件,通过分区机制实现水平扩展,副本策略保障数据高可用;
  • 实时计算层:Flink 负责 ETL 清洗、维度关联与高频指标聚合,基于事件时间与 Watermark 机制保障数据一致性;
  • 存储查询层:ClickHouse 负责数据持久化与复杂指标预计算,通过物化视图提升查询效率;
  • 应用层:支撑实时报表、佣金结算、全链路监控等核心业务场景。

2. 核心技术选型依据

技术组件

选型结果

备选方案

决策核心依据

实时计算引擎

Flink

Spark Streaming

Exactly-once 语义保障数据一致性,毫秒级延迟满足实时需求,丰富状态管理适配复杂计算

OLAP 分析引擎

ClickHouse

Druid

列式存储适配宽表聚合,物化视图预计算能力突出,分布式架构支撑 5000+QPS 高并发查询

消息队列

Kafka

-

高吞吐量支撑海量数据传输,动态分区机制适配业务增长,副本机制保障数据不丢失

缓存组件

Redis

-

支撑静态维度数据缓存,降低数据库查询压力,提升维度关联效率

三、分层架构核心:DWD+DWS 设计与实现

        分层设计是保障架构灵活性与复用性的关键,通过 DWD 层构建标准化数据底座,DWS 层实现双模式指标聚合,最终达成 “数据一次处理、多端复用” 的目标。

1. DWD 层(明细数据层):标准化数据底座

核心定位

        对原始数据进行轻量级清洗、格式标准化与维度关联,生成高质量明细宽表,解决 “数据口径不统一、重复清洗、维度缺失” 问题,为下游计算提供统一数据来源。

核心处理逻辑
  • 数据清洗:过滤无效数据(如取消订单、空值数据)、去重处理、补全缺失字段(如默认佣金比例);
  • 维度关联:
    • 静态维度(商品品类、推广渠道):采用 Flink Lookup Join+Redis 缓存,预加载维度数据,避免频繁访问数据库;
    • 动态维度(用户会员等级、实时汇率):采用 Flink Interval Join(双流 Join),设置 5 分钟时间窗口,结合 Watermark 处理数据乱序;
  • 格式标准化:统一字段命名、数据类型(如金额统一为 Decimal (18,2)),保障数据一致性。
存储策略

        双落地策略:写入 Kafka DWD 层 Topic(支撑下游 Flink 实时聚合),同时写入 ClickHouse DWD 明细底表(支撑复杂查询与物化视图预计算),支持数据回溯与重计算。

扩展性设计

        预留维度扩展字段,采用 “3-5 个固定字段 + JSON 嵌套字段” 组合方案,应对业务快速迭代:

  • 固定字段:ext_dimension_1/2/3,适配明确的扩展需求;
  • JSON 字段:ext_info,内部键名采用小驼峰命名(如 memberLevel、productTags),支撑不确定维度的灵活扩展;
  • 新增业务维度时,仅需补充 DWD 层关联逻辑,下游无需修改。

2. DWS 层(指标服务层):双模式指标聚合

核心定位

        基于 DWD 层明细数据,实现 “高频实时指标 + 低频复杂指标” 双模式聚合,封装通用计算逻辑,提供标准化数据服务接口。

双模式聚合实现
  • 实时聚合模式(Flink 主导):
    • 适用场景:高频、低延迟需求的指标(如实时订单数、佣金总额、渠道转化率);
    • 实现方式:Flink 滚动窗口(1 分钟)聚合,结果直接写入 ClickHouse DWS 指标表,支撑实时报表与结算业务;
  • 预计算模式(ClickHouse 物化视图主导):
    • 适用场景:低频、复杂查询需求的指标(如 “渠道 + 品类 + 会员等级” 交叉分析);
    • 实现方式:基于 ClickHouse DWD 明细底表创建物化视图,预计算并存储聚合结果,避免查询时扫描全量明细,将查询耗时从分钟级降至毫秒级。
存储与服务形式
  • 存储:实时聚合结果写入 ClickHouse DWS 指标表 + Kafka DWS Topic;预计算结果存储为按业务维度拆分的物化视图(如 mv_commission_channel_category);
  • 服务接口:
    • RESTful API:部署于微服务网关,统一返回格式(code/data/msg),支持 Swagger 在线调试;
    • 数据订阅:支持 Kafka Topic 订阅与 ClickHouse JDBC 直连,满足不同业务场景需求;
    • 权限控制:集成 OAuth2.0 统一鉴权,按角色控制数据可见范围。

四、关键技术实现与性能优化

1. 维度关联优化:静态与动态维度差异化处理

  • 静态维度(更新慢、数据量小):Redis 缓存 + Flink Lookup Join,查询响应时间降至毫秒级,降低数据库压力;
  • 动态维度(更新快、数据量中):Kafka 双流 Join+Watermark 机制,设置 5 分钟时间窗口匹配最新数据,引入维度版本号优先选择最新记录,解决数据乱序问题,佣金计算准确率提升至 99.99%。

2. ClickHouse 性能优化核心策略

  • 物化视图优化:按业务维度拆分视图,采用 “日志表 + MergeTree” 架构实现增量刷新,避免全量刷新损耗;高频视图 5 分钟刷新,低频视图 1 小时刷新,适配不同查询场景;
  • 存储优化:按 “时间 + 业务维度” 分片,按天分区存储,DWD 明细底表设置 TTL 自动清理 90 天前数据,降低存储压力;
  • 索引优化:为高频查询的物化视图建立稀疏索引,提升扫描效率,复杂查询耗时降低 85%。

3. Flink 性能与可靠性优化

  • 窗口优化:采用 1 分钟滚动窗口平衡实时性与计算压力,优化 Watermark 生成策略,减少数据乱序影响;
  • 状态管理:使用 RocksDB 状态后端,支持大状态高效读写,通过 Checkpoint/Savepoint 机制保障故障恢复能力,实现端到端 Exactly-once 语义;
  • 资源调度:基于 YARN + 资源调度平台,支持 TaskManager 弹性扩缩容,高峰时段从 10 台扩容至 50 台,适配数据量波动。

4. 全链路监控告警体系

  • 告警场景:覆盖数据延迟、指标异常、系统故障三类核心场景;
  • 触发逻辑:
    • 数据延迟:Flink Watermark 监控,延迟超 30 秒直接告警;
    • 指标异常:Flink 聚合时嵌入判断逻辑,连续 2 个窗口指标超标(如佣金误差率 > 5%)触发告警;
    • 系统故障:SkyWalking+Prometheus 采集 CPU 利用率、查询耗时等指标,按阈值触发告警;
  • 保障机制:支持业务方可视化配置阈值,设置 5 分钟告警抑制,避免告警风暴,关联原始数据 ID 便于快速排查。

5. 数据治理与血缘设计(支撑架构可靠性)

数据治理与血缘是保障架构长期稳定运行的核心支撑,通过体系化管理实现数据 “高质量、可追溯、安全合规”。

(1)数据治理落地
  • 数据质量保障:DWD 层规则引擎过滤无效数据、补全缺失字段(清洗覆盖率 100%);Flink 聚合时校验指标合理性(如佣金 > 0),异常数据单独归档;每小时比对 DWD 明细与 DWS 聚合总数,差异超 0.1% 触发告警。
  • 数据安全管控:按 “所有者 - 管理者 - 使用者” 分级授权,推广人员仅见负责渠道数据;DWD 层脱敏用户手机号、身份证号(如 138****5678),规避泄露风险。
  • 数据生命周期管理:DWD 明细存 90 天(ClickHouse),DWS 指标存 1 年(ClickHouse + 低成本存储);超 3 个月数据迁移至冷存储,ClickHouse 仅留近 3 个月热数据,存储成本降 70%。
(2)数据血缘追踪
  • 血缘采集:基于 Flink 日志记录 “数据源→DWD→DWS” 流转关系及指标计算逻辑(如commission_amount=order_amount*commission_rate);基于 ClickHouse 表关联,采集物化视图与源表依赖关系。
  • 价值与应用:可视化展示全链路血缘,支持正向追溯(数据源→下游指标)与反向追溯(指标→数据源);指标异常时快速定位问题数据源(故障排查从 2 小时缩至 10 分钟);维度变更前分析下游影响,避免业务中断。

五、项目落地效果与技术亮点

1. 核心指标提升

  • 数据延迟:从小时级降至 200ms 内,满足实时结算需求;
  • 查询性能:复杂报表查询耗时降低 85%,支撑 5000+QPS 高并发查询;
  • 扩展性:新增 3 个业务线仅需 1 周适配,无需重构核心代码;
  • 数据准确性:佣金核算准确率达 99.99%,核心业务链路零故障。

2. 核心技术亮点

  • 架构创新:采用 “Flink 实时聚合 + ClickHouse 物化视图预计算” 双模式,平衡实时性与查询性能,解决传统架构 “实时性与性能不可兼得” 的痛点;
  • 分层解耦:DWD+DWS 分层设计实现数据 “一次清洗、多端复用”,减少重复计算,数据处理效率提升 60%;
  • 工程化保障:通过全链路监控、数据治理与血缘追踪,故障定位时间从 2 小时缩短至 10 分钟,架构可靠性达 99.99%;
  • 复用性强:标准化接口与扩展字段设计,支持多业务线快速接入,新增业务线适配周期从 1 周缩短至 1 天。

总结

        千万级订单实时计算场景的核心是 “分层解耦 + 技术协同”,DWD 层构建标准化数据底座保障数据质量与扩展性,DWS 层双模式聚合平衡实时性与查询性能,Flink 与 ClickHouse 的协同则最大化发挥各自技术优势。

        本文分享的架构设计与优化思路已在实际项目中验证,可广泛应用于电商推广、金融结算、实时监控等同类场景,为大数据实时处理提供可复用的实践参考。

若有架构设计细节、Flink 状态管理或 ClickHouse 优化相关问题,欢迎在评论区交流讨论!

关联阅读

​​​​​​​


📚 我的技术博客导航:[点击进入一站式查看所有干货]


Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐