实时数据看板的架构重构实践:从数据库轮询到 Flink + Redis 流式计算
实时数据看板的架构重构实践从数据库轮询到 Flink Redis 流式计算一、当仪表盘遇上海量数据数据库直查模式的全面溃败任何一个面向运营或管理层的数据看板都会经历从页面秒开到点一下卡三秒的退化过程。根本原因在于随着业务增长维度的组合爆炸与数据的实时性要求形成了双重压力。某电商平台的核心运营看板——涵盖 GMV、订单量、转化率、退货率等 20 指标按小时、商品类目、地域做多维度下钻——在日均订单突破 50 万后页面加载时间从 200ms 恶化到 8 秒以上。问题的根源有两个层面。首先是查询层面每个看板组件的刷新都对应一条或多条聚合 SQL。一条典型的多维聚合查询——按省份、类目、时段分组的订单汇总——在千万级订单表上执行 GROUP BY即使建立了覆盖索引执行时间也在 3-5 秒。看板上有 12 个组件同时刷新即使做了并发查询总耗时也远超用户容忍度。其次是资源层面高峰时期 OLTP 数据库同时承载交易写入和看板查询两者互相争抢 CPU 和磁盘 IO。分析发现看板查询贡献了约 35% 的数据库 CPU 消耗而其中 60% 的查询结果在秒级之内并无变化——这是典型的用 OLTP 资源执行 OLAP 任务的资源错配。解决方向很明确将看板的数据计算从 OLTP 中剥离出来迁入流式计算管道用预聚合替代实时查询用缓存替代数据库直读。二、流批一体的看板数据管道Flink Kafka Redis 三层架构重构后的架构将数据链路拆分为三个清晰的层次数据采集层使用 Canal 监听 MySQL Binlog将订单、支付、退款等核心业务的变更事件实时投递到 Kafka。同时前端埋点数据页面浏览、商品点击等也汇入同一个 Kafka 集群。这种设计保证了业务数据 行为数据在同一套管道中流转。流式计算层的核心是 Flink。这里选择了多窗口设计1 分钟滚动窗口覆盖实时大屏的高频刷新需求5 分钟滑动窗口覆盖运营看板的准实时需求1 小时滚动窗口的数据则写入 ClickHouse 用于历史趋势分析。三种窗口的计算结果独立存储互不干扰。窗口计算之后是维表 Join。以商品类目为例订单流中只有商品 ID需要关联类目信息才能做按类目的聚合。维表数据量不大几千到几万条加载到 Flink 的 RocksDB State Backend 中做本地 Join避免了每条数据都查一次 Redis 的网络开销。存储输出层采用 Redis ClickHouse 的异构存储方案。Redis 存储热数据——最近 1 小时内的分钟级和 5 分钟级聚合结果TTL 设为 2 小时自动清理。ClickHouse 存储按小时的聚合结果用于历史趋势查询和分析报表。在 API 服务端看板请求优先读 Redis。如果 Redis 中对应时间窗口的数据不存在冷启动或数据过期则兜底查询 ClickHouse并将结果回写到 Redis 中Cache-Aside 模式。三、Flink 多窗口聚合与 Redis 原子写入的核心实现以下是 Flink 作业中多窗口聚合与结果写入的关键代码public class DashboardAggregationJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment .getExecutionEnvironment(); // 开启 Checkpoint精确一次语义保障 env.enableCheckpointing(60000); env.getCheckpointConfig().setCheckpointingMode( CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setCheckpointTimeout(120000); // 数据源Kafka 订单流 DataStreamOrderEvent orderStream env .addSource(new FlinkKafkaConsumer( order-events, new OrderEventDeserializationSchema(), KafkaConfig.buildProperties())) .assignTimestampsAndWatermarks( WatermarkStrategy.OrderEventforBoundedOutOfOrderness( Duration.ofSeconds(10)) .withTimestampAssigner((event, ts) - event.getTimestamp())) .name(kafka-order-source); // 1 分钟滚动窗口聚合 DataStreamDashboardMetric minuteAgg orderStream .keyBy(OrderEvent::getRegionKey) // 按地域类目组合键分区 .window(TumblingEventTimeWindows.of(Time.minutes(1))) .aggregate(new OrderAggregateFunction()) .name(minute-window-agg); // 5 分钟滑动窗口聚合 DataStreamDashboardMetric slideAgg orderStream .keyBy(OrderEvent::getRegionKey) .window(SlidingEventTimeWindows.of( Time.minutes(5), Time.minutes(1))) .aggregate(new OrderAggregateFunction()) .name(minute5-slide-window-agg); // 写入 Redis热数据 minuteAgg.addSink(new RedisDashboardSink(minute)); slideAgg.addSink(new RedisDashboardSink(minute5)); env.execute(Dashboard Real-time Aggregation Job); } /** * 自定义聚合函数在窗口内计算 GMV、订单数、退款数和退款率 */ public static class OrderAggregateFunction implements AggregateFunctionOrderEvent, MetricAccumulator, DashboardMetric { Override public MetricAccumulator createAccumulator() { return new MetricAccumulator(); } Override public MetricAccumulator add(OrderEvent event, MetricAccumulator acc) { switch (event.getType()) { case ORDER_CREATED: acc.addOrder(event.getAmount()); break; case ORDER_PAID: acc.addPaid(event.getAmount()); break; case ORDER_REFUNDED: acc.addRefund(event.getAmount()); break; // 忽略其他事件类型避免脏数据干扰聚合 } return acc; } Override public DashboardMetric getResult(MetricAccumulator acc) { double refundRate acc.getOrderCount() 0 ? (double) acc.getRefundCount() / acc.getOrderCount() : 0.0; return new DashboardMetric( acc.getGmv(), acc.getOrderCount(), acc.getPaidAmount(), acc.getRefundAmount(), // 保留 4 位小数避免浮点精度问题 Math.round(refundRate * 10000.0) / 10000.0 ); } Override public MetricAccumulator merge(MetricAccumulator a, MetricAccumulator b) { a.merge(b); return a; } } }Redis 写入 Sink 的设计有一个容易被忽略的细节原子性。一个时间窗口内可能有数十个维度的聚合结果需要写入 Redis如果逐条 SET在写入过程中发生故障会导致部分数据缺失。我们的做法是使用 Redis Pipeline 做批量写入配合 Flink 的 Checkpoint 做事务性提交public class RedisDashboardSink extends RichSinkFunctionDashboardMetric { private transient JedisCluster jedisCluster; Override public void open(Configuration parameters) { jedisCluster new JedisCluster( new HostAndPort(redis-cluster, 6379), 2000, 2000, 3, password, new GenericObjectPoolConfig()); } Override public void invoke(DashboardMetric metric, Context context) { String key buildRedisKey(metric); String value JSON.toJSONString(metric); // 使用 SETEX 原子操作写入数据同时设置过期时间 // TTL 根据窗口粒度差异化分钟级 2h5 分钟级 4h int ttlSeconds windowType.equals(minute) ? 7200 : 14400; try { jedisCluster.setex(key, ttlSeconds, value); } catch (JedisConnectionException e) { log.error(Redis 写入失败key{}, 将在下次 Checkpoint 重试, key, e); // 抛出异常触发 Flink 重试不影响数据一致性 throw new RuntimeException(Redis write failed for key: key, e); } } }这里的关键决策是选择抛出异常而非吞掉异常。如果吞掉异常Flink 不会感知到写入失败该条数据将丢失抛出异常后Flink 会根据 Checkpoint 机制回滚到上一个成功位点并自动重试。代价是重试期间可能产生重复数据但因聚合结果本身是幂等的同 key 覆盖写入这并不影响最终一致性。四、状态膨胀与延迟抖动流式方案并非万能药这个方案有几个明确的边界条件必须认清。首先Flink 的状态管理开销。随着聚合维度的增加——比如将维度从 10 个省份扩展到 2000 个区县——RocksDB State 的大小会迅速膨胀。每个 key 的窗口聚合状态约 200 字节十万级 key 的 State 总量约 20MB还处于可控范围但如果维度进一步细化到百万级就需要考虑使用增量 Checkpoint 和 State TTL 来控制状态大小。其次Redis 的内存成本不可忽视。以 1 分钟粒度为例如果看板有 100 个维度组合每个聚合结果的 JSON 约 500 字节保留 2 小时就是 60 × 100 × 500 ≈ 3MB——这还好。但生产环境中维度组合往往远超 100 个实际可能达到数万个内存占用会在数十 GB 级别。必须严格控制 Redis 中存储的维度范围只缓存高频查询的维度组合。第三个容易被忽视的问题是数据延迟的累积效应。Canal Binlog 解析约 50msKafka 传输约 20msFlink 窗口触发需要等待 Watermark 推进至少 10 秒我们配置的乱序容忍时间。这意味着端到端延迟从数据产生到看板刷新大约在 12-15 秒之间。这个延迟对大屏展示可以接受但对交易风控等场景远远不够——那些场景需要的是 CEP 模式匹配而非窗口聚合。五、总结从数据库直查到流式计算的看板架构重构本质上是一次计算前置化和存储特化的改造。将聚合计算从查询时转义到数据流入时执行用预计算换取低延迟将存储从通用 OLTP 存储迁移到 Redis ClickHouse 异构存储用专用化换高性能。落地节奏建议第一阶段先将流量最大、延迟敏感的 3-5 个核心指标接入 Flink 管道验证端到端延迟和数据准确性第二阶段扩覆盖到全量指标建立 Redis 热数据 ClickHouse 冷数据的双轨存储第三阶段引入 Flink CDC 替代 Canal简化数据采集层的链路复杂度。监控层面重点关注 Flink Checkpoint 时长、Kafka Consumer Lag、Redis 内存使用率和看板 API 的 P99 延迟这四项指标。