实时数仓是现在大数据领域的热门方向。

以前我们做数据分析,都是T+1的,今天看昨天的数据。但现在越来越多的场景需要实时数据:实时大屏看销售数据、实时推荐商品、实时风控识别欺诈、实时监控系统状态。

我最近参与了公司实时数仓的搭建,从0到1走了一遍,踩了不少坑。今天分享实时数仓的入门知识,包括什么是实时数仓、核心组件、架构设计、具体实现步骤,以及踩过的坑。帮你快速入门实时数仓。

一、什么是实时数仓

在说实时数仓之前,先说说传统数仓。

传统数仓(离线数仓)是T+1的。每天晚上,把业务数据库的数据同步到数仓,经过清洗、转换、聚合,第二天早上就能看到前一天的数据报表。

这种方式的优点是稳定、成本低,缺点是延迟高,只能看历史数据,不能实时反映业务状况。

实时数仓,就是把数据处理的延迟从"天"降到"秒"甚至"毫秒"。数据产生后,经过实时处理,几秒内就能查询到结果。

比如,用户下单后,实时数仓几秒内就能把这笔订单统计到销售数据里,大屏上的数字立刻更新。

实时数仓的核心是"实时",但它的本质还是数仓,要遵循数仓的分层设计、数据质量、数据治理等原则。

二、为什么需要实时数仓

实时数仓不是银弹,不是所有场景都需要。但以下场景,实时数仓很有价值。

1. 实时大屏

老板想看实时的销售数据、订单量、用户活跃度,大屏上的数字要实时更新。这就需要实时数仓。

2. 实时推荐

电商、视频、新闻等平台,需要根据用户的实时行为推荐内容。用户刚看了什么、买了什么,要立刻反映到推荐结果里。

3. 实时风控

金融、电商等场景,需要实时识别欺诈行为。比如用户下单后,要立刻判断是不是欺诈订单,是就拦截,不是就放行。

4. 实时监控

系统监控、业务监控,需要实时发现异常。比如服务器CPU突然升高、订单量突然下降,要立刻告警。

5. 实时报表

运营人员需要看实时的报表,了解当前的业务状况,及时做出决策。

如果你的业务有这些需求,就可以考虑搭建实时数仓。

三、实时数仓的核心组件

搭建实时数仓,需要几个核心组件。

1. 数据采集:Canal / Flume / CDC

实时数仓的第一步,是采集数据。

业务数据存在MySQL、PostgreSQL等关系型数据库里。要实时采集这些数据,可以用Canal、Debezium等工具,监听数据库的binlog,把数据变更实时同步出来。

日志数据(用户行为日志、应用日志)可以用Flume、Logstash、Filebeat等工具采集。

采集到的数据,统一发到消息队列里。

2. 消息队列:Kafka

Kafka是实时数仓的核心组件,用来做数据缓冲和消息队列。

为什么需要Kafka?因为数据采集的速度和数据处理的速度不一定匹配。比如大促的时候,订单量暴增,处理能力跟不上,如果没有缓冲,数据就会丢。

Kafka作为缓冲,把采集到的数据先存起来,处理引擎慢慢消费。这样即使处理慢一点,也不会丢数据。

Kafka还有一个好处是解耦。采集端和处理端互不影响,各自独立扩展。

3. 实时计算:Flink

Flink是目前最主流的实时计算引擎,也是实时数仓的核心。

Flink能从Kafka消费数据,进行清洗、转换、聚合、关联等操作,然后把结果写到下游存储。

Flink的特点:

  • 低延迟:毫秒级延迟
  • 高吞吐:每秒能处理几百万条数据
  • 精确一次(Exactly-Once):保证数据不丢不重
  • 支持事件时间:能处理乱序数据
  • 有状态计算:支持复杂的聚合和关联

除了Flink,还有Spark Streaming、Storm等,但目前Flink是实时计算的首选。

4. 数据存储:ClickHouse / Doris / HBase

实时计算的结果,要存到适合实时查询的存储里。

常用的实时数仓存储:

  • ClickHouse:列式存储,查询速度极快,适合OLAP分析。单表查询性能很强。
  • Apache Doris:MPP架构,支持高并发查询,适合实时报表和大屏。
  • HBase:KV存储,适合点查,不适合复杂分析。
  • Redis:内存数据库,速度快,适合存简单的指标,但不适合复杂查询。

根据查询场景选择合适的存储。如果是复杂分析查询,选ClickHouse或Doris;如果是简单的点查,选HBase或Redis。

5. 查询服务:Presto / 自定义API

数据存好了,还要提供查询服务。

可以用Presto、Trino等查询引擎,直接查询ClickHouse、Hive等存储。也可以自己写API,封装查询逻辑,提供给前端或其他系统调用。

6. 数据可视化:Superset / Grafana

最后,用可视化工具把数据展示出来。

  • Superset:开源BI工具,支持多种数据源,能做各种图表和大屏。
  • Grafana:适合做监控大屏,和时序数据库配合好。
  • FineBI / Tableau:商业BI工具,功能更强。

四、实时数仓的架构设计

实时数仓的架构,一般分为几层。

数据源 → 数据采集 → 消息队列 → 实时计算 → 数据存储 → 查询服务 → 可视化

具体来说:

1. ODS层(原始数据层)

ODS层存原始数据,和数据源保持一致,不做清洗和转换。

实时数仓的ODS层,一般就是Kafka里的原始数据。Canal采集的binlog、Flume采集的日志,都先写到Kafka的ODS topic里。

2. DWD层(明细数据层)

DWD层是清洗后的明细数据。从ODS层消费数据,做清洗(去重、过滤、格式转换)、维度关联,然后写到Kafka的DWD topic里,或者直接写到存储里。

DWD层的数据是干净的、结构化的明细数据,是后续计算的基础。

3. DWS层(汇总数据层)

DWS层是轻度汇总的数据。按主题(比如订单、用户、商品)做轻度聚合,比如按天、按小时汇总。

DWS层的数据存在ClickHouse或Doris里,供查询使用。

4. ADS层(应用数据层)

ADS层是面向应用的数据。根据具体的业务需求,做进一步的聚合和计算,比如实时大屏的数据、推荐的特征、风控的指标。

ADS层的数据直接供应用查询。

这种分层设计,和离线数仓类似,好处是解耦、可复用、便于维护。

五、从0到1搭建步骤

下面说说具体怎么从0到1搭建一个简单的实时数仓。

第一步:确定需求和数据源

先想清楚,你要做什么实时分析?需要哪些数据?

比如,要做一个实时销售大屏,需要订单数据、商品数据、用户数据。这些数据存在哪些数据库里?哪些表?

需求明确了,才知道要采集什么数据。

第二步:搭建Kafka集群

安装部署Kafka,创建需要的topic。

比如:

  • ods_order:订单原始数据
  • ods_user:用户原始数据
  • dwdorderdetail:清洗后的订单明细
  • dwduserdetail:清洗后的用户明细

Kafka集群至少3个节点,保证高可用。

第三步:配置数据采集

用Canal监听MySQL的binlog,把订单表、用户表的变更实时同步到Kafka的ODS topic里。

Canal的配置很简单,指定要监听的数据库和表,以及目标Kafka的topic就行。

日志数据用Flume或Filebeat采集,也写到Kafka里。

第四步:开发Flink实时计算任务

这是最核心的一步。用Flink开发实时计算任务:

  1. 从Kafka的ODS topic消费数据
  2. 做数据清洗:去重、过滤脏数据、格式转换
  3. 做维度关联:把订单数据和商品维度、用户维度关联起来
  4. 做聚合计算:按小时、按商品类别汇总销售额、订单量
  5. 把结果写到ClickHouse或Doris里

Flink任务可以用Java/Scala开发,也可以用Flink SQL。Flink SQL更简单,适合做ETL和简单聚合。

举个简单的Flink SQL例子:

-- 从Kafka读取订单数据
CREATE TABLE order_ods (
  order_id BIGINT,
  user_id BIGINT,
  product_id BIGINT,
  amount DECIMAL(10,2),
  create_time TIMESTAMP,
  WATERMARK FOR create_time AS create_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'ods_order',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'json'
);

-- 按小时汇总销售额
CREATE TABLE hourly_sales (
  hour TIMESTAMP,
  total_amount DECIMAL(10,2),
  order_count BIGINT
) WITH (
  'connector' = 'clickhouse',
  'url' = 'clickhouse:8123',
  'table-name' = 'hourly_sales',
  'database-name' = 'dws'
);

INSERT INTO hourly_sales
SELECT 
  TUMBLE_START(create_time, INTERVAL '1' HOUR) as hour,
  SUM(amount) as total_amount,
  COUNT(*) as order_count
FROM order_ods
GROUP BY TUMBLE(create_time, INTERVAL '1' HOUR);

这个例子就是从Kafka读订单数据,按小时汇总销售额和订单量,写到ClickHouse里。

第五步:搭建数据存储

安装部署ClickHouse或Doris,创建需要的表。

表设计要注意:

  • 按时间分区,方便管理和查询
  • 选择合适的排序键,提升查询性能
  • 合理设置TTL,自动清理过期数据

第六步:开发查询服务和可视化

用Superset或Grafana连接ClickHouse,创建实时大屏和报表。

或者自己写API,封装查询逻辑,提供给前端调用。

第七步:监控和运维

实时数仓是7x24小时运行的,监控很重要。

要监控:

  • Kafka的消息积压情况
  • Flink任务的运行状态、延迟、吞吐量
  • ClickHouse的查询性能、存储用量
  • 数据质量:数据是否丢失、是否准确

监控告警做好了,才能及时发现问题,保证实时数仓的稳定运行。

六、踩过的坑

搭建实时数仓的过程中,踩了不少坑,分享一下。

坑一:数据乱序

实时数据经常是乱序的。比如订单10:00创建,10:05才同步到Kafka;或者网络延迟,后发生的事件先到。

Flink用Watermark机制处理乱序数据,但Watermark的设置很关键。设置太短,乱序数据会被丢弃;设置太长,延迟增加。

建议根据业务场景,合理设置Watermark,一般允许5-30秒的乱序。

坑二:数据重复

实时处理中,数据重复是常见问题。比如Flink任务重启,可能会重复处理数据。

Flink的Exactly-Once语义能保证不丢不重,但需要正确配置Checkpoint,并且下游存储支持幂等写入。

ClickHouse的ReplacingMergeTree能去重,但查询的时候要注意。

坑三:大状态

Flink做聚合和关联的时候,会产生状态。如果数据量大,状态会很大,影响性能,甚至导致OOM。

解决方法:

  • 合理设置状态后端(RocksDB)
  • 设置状态TTL,自动清理过期状态
  • 优化计算逻辑,减少不必要的状态

坑四:维度关联的时效性

实时数仓经常需要关联维度表(比如商品信息、用户信息)。但维度表是会变化的,怎么保证关联的是正确的维度?

常用的方法:

  • 用Flink的广播流,把维度表广播到所有并行实例
  • 用Lookup Join,实时查询维度表(但性能差)
  • 用CDC同步维度表的变更,实时更新

坑五:实时和离线的数据不一致

很多公司同时有离线数仓和实时数仓,两者的数据经常对不上。

原因可能是:

  • 计算逻辑不一致
  • 数据口径不一致
  • 实时数据有延迟或丢失

建议:统一指标口径,定期用离线数据校准实时数据。

坑六:成本控制

实时数仓的成本不低。Flink集群、Kafka集群、ClickHouse集群,都需要不少服务器。

要注意成本控制:

  • 合理设置并行度,不要浪费资源
  • 冷热数据分离,热数据存ClickHouse,冷数据存HDFS
  • 自动扩缩容,高峰期扩容,低峰期缩容

七、学习建议

如果你想学习实时数仓,给几个建议。

1. 先学基础

先了解大数据的基础知识:Hadoop、Hive、Kafka、Flink。不用精通,但要知道每个组件是做什么的。

2. 动手实践

光看没用,要动手。在自己的电脑上搭个伪分布式,跑一个简单的实时计算任务。

比如,用Canal采集MySQL的数据,写到Kafka,用Flink做简单聚合,写到ClickHouse,用Superset展示。跑通这个流程,就入门了。

3. 看官方文档

Flink、Kafka、ClickHouse的官方文档都很详细,是最好的学习资料。

4. 关注社区

实时数仓发展很快,新技术、新框架不断出现。关注社区,了解最新动态。

5. 从简单开始

不要一开始就搞很复杂的架构。先做一个简单的实时数仓,跑通流程,再慢慢优化和扩展。

八、写在最后

实时数仓是一个很有前景的方向。随着业务对实时性的要求越来越高,实时数仓会越来越普及。

但实时数仓也不是银弹。它成本高、复杂度高、运维难。如果你的业务不需要实时数据,就没必要上实时数仓,离线数仓足够了。

如果确实需要实时数据,那就从简单开始,一步步搭建。先跑通核心流程,再慢慢完善。

2022年了,实时计算的技术越来越成熟,Flink、ClickHouse这些工具也越来越好用。搭建实时数仓的门槛在降低,中小公司也能用上实时数仓了。

希望这篇文章能帮你入门实时数仓。如果你有问题或经验,欢迎在评论区交流。

最后,用一句话总结:"实时数仓的核心不是实时,而是数仓。先把数仓的基础打好,再谈实时。"

祝大家都能搭建出稳定、高效的实时数仓。