尚硅谷Flink实时数仓3.0抢先版:手把手教你构建电商实时分析系统
从零到一:基于Flink 1.17+的电商实时数仓3.0实战全解析
最近和几个在一线大厂做数据平台的朋友聊天,大家不约而同地提到了一个词:实时化焦虑。无论是业务方对“秒级”数据看板的催促,还是风控对“毫秒级”异常行为拦截的需求,都让传统的T+1离线数仓显得力不从心。这背后,其实是整个行业从“事后分析”向“事中决策”甚至“事前预测”的范式转移。如果你正面临类似的挑战,或者想系统性地掌握构建一套高可用、高性能实时数据系统的能力,那么这篇文章或许能给你带来一些不一样的思路。我们不谈空洞的理论,而是聚焦于一个具体的电商场景,手把手拆解如何用当前最主流的Flink生态,搭建一套代号为“3.0”的实时数仓。这套架构不仅融合了最新的Flink CDC、流批一体和实时OLAP技术,更在数据一致性、开发效率和运维成本之间找到了新的平衡点。
1. 架构演进与3.0核心设计理念
在深入代码之前,我们有必要先理解为什么是“3.0”。回顾实时数仓的演进,大致可以划分为几个阶段:
- 1.0时代(Lambda架构):这是早期的经典模式,一条实时流处理管道(如Storm/Spark Streaming)提供低延迟但可能不精确的结果,另一条离线批处理管道(如Hive/Spark)在次日提供精准的“修正”结果。开发维护两套逻辑,数据一致性挑战巨大。
- 2.0时代(Kappa架构与流批一体初步):尝试用一套流处理系统(如Flink)解决所有问题,通过将历史数据重新灌入流处理系统来模拟批处理。这简化了架构,但对消息队列(如Kafka)的长期存储能力和流处理引擎的吞吐量提出了极高要求。
- 3.0时代(流批一体与数据湖仓融合):这是当前的前沿实践。其核心思想不再是“用流模拟批”或反之,而是在计算层实现API的统一,在存储层实现流批数据的统一存储与管理。Flink的流批一体SQL引擎,配合Iceberg、Hudi这类数据湖表格式,使得同一份数据既能支持低延迟的流式读取,也能支持高性能的批式回溯分析。
我们本次构建的电商实时数仓3.0,正是基于这一理念。它的核心目标有三个:端到端秒级延迟、 Exactly-Once语义保障、以及面向分析师的高效即席查询能力。整个架构可以概括为下图(文字描述):
注意:与2.0架构重度依赖Kafka作为中间层不同,3.0架构中,Kafka的角色更偏向于高性能、解耦的实时数据管道,而将具备事务能力的数据湖表(如Apache Iceberg) 作为流批数据的统一存储层和批处理任务的源表。
为了更清晰地对比不同架构的特点,我们可以参考下表:
| 架构版本 | 核心特征 | 存储层 | 计算层 | 优点 | 挑战 |
|---|---|---|---|---|---|
| Lambda (1.0) | 实时流+离线批双路并行 | Kafka + HDFS | Storm/Spark Streaming + Hive/Spark | 技术栈成熟,实时离线隔离 | 逻辑重复,数据不一致,运维复杂 |
| Kappa (2.0) | 一切皆流,重放历史 | Kafka(长期存储) | Flink | 架构简化,一套逻辑 | Kafka存储压力大,历史数据分析成本高 |
| 流批一体 (3.0) | 统一存储,统一API | Kafka + Iceberg/Hudi | Flink (流批一体SQL) | 存储计算解耦,数据一致性高,支持高效批查询 | 技术栈较新,生态工具链仍在完善 |
我们的3.0架构设计,简单来说就是:使用Flink CDC直接捕获MySQL的变更日志,写入Kafka作为实时数据流;Flink流作业消费Kafka,进行实时ETL和聚合,同时将结果实时写入Iceberg表;最终,通过Doris或ClickHouse这类实时OLAP数据库对Iceberg表进行加速查询,供BI工具或数据应用使用。这样,实时流处理、离线数据补全、交互式分析都基于同一份Iceberg数据,彻底解决了数据孤岛问题。
2. 环境搭建与核心组件选型
工欲善其事,必先利其器。搭建一套可实战的实时数仓环境,组件的选型和版本搭配至关重要。这里我推荐一套经过生产环境验证的稳定组合,你完全可以在自己的虚拟机或云服务器上复现。
基础环境准备:
- 操作系统:CentOS 7.9 或 Ubuntu 20.04 LTS。
- Java:JDK 11(Flink 1.17+的官方推荐版本,务必使用Oracle JDK或OpenJDK,避免使用有兼容性问题的其他发行版)。
- 集群资源:建议至少3个节点,配置4核8G以上,保证基本的冗余和高可用。
核心组件清单与版本:
- Apache Flink 1.17.1:我们选择的计算引擎核心。1.17版本在流批一体SQL的成熟度、与数据湖集成的易用性上都有显著提升。
- Apache Kafka 3.4.0:作为实时数据总线。选择3.x版本是为了更好的性能和对Kraft模式(去ZooKeeper)的支持,简化部署。
- Apache Iceberg 1.3.0:作为统一的数据湖存储格式。它提供了ACID事务、隐藏分区、模式演化等企业级特性,与Flink集成度极高。
- Flink CDC 2.4+:用于捕获MySQL变更数据的神器。我们使用
flink-sql-connector-mysql-cdc,它能直接以SQL的方式定义CDC源表,极大简化开发。 - 实时OLAP引擎(二选一):
- Apache Doris 1.2.4:国产优秀项目,对宽表聚合查询和标准SQL支持非常好,运维相对简单。
- ClickHouse 22.8+:在复杂聚合查询和单表性能上表现极致,但分布式Join和更新操作是其弱项。
- 监控与调度:Prometheus + Grafana 监控Flink和Kafka;Apache DolphinScheduler 用于调度离线补数据任务。
这里给出一个使用Docker Compose快速拉起Flink单机版、Kafka和MySQL(作为数据源)的示例,用于本地开发测试:
# docker-compose.yml
version: '2.1'
services:
jobmanager:
image: flink:1.17.1-scala_2.12-java11
ports:
- "8081:8081"
command: jobmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
taskmanager.numberOfTaskSlots: 4
taskmanager:
image: flink:1.17.1-scala_2.12-java11
depends_on:
- jobmanager
command: taskmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
taskmanager.numberOfTaskSlots: 4
parallelism.default: 2
mysql:
image: debezium/mysql:8.0
ports:
- "3306:3306"
environment:
- MYSQL_ROOT_PASSWORD=123456
- MYSQL_USER=flinkuser
- MYSQL_PASSWORD=flinkpw
- MYSQL_DATABASE=ecommerce
zookeeper:
image: confluentinc/cp-zookeeper:7.3.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: confluentinc/cp-kafka:7.3.0
depends_on:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
使用 docker-compose up -d 启动后,你就拥有了一个最小化的实时数据处理环境。生产环境则需要分布式部署,并仔细配置网络、存储和各项参数。
3. 数据接入层:Flink CDC与实时数据流构建
数据接入是实时数仓的“水源”。传统上,我们会用Canal或Debezium监听MySQL binlog,然后编写程序解析并写入Kafka。现在,Flink CDC让这一切变得异常简单。它直接将MySQL的binlog变更事件作为Flink SQL中的一个动态表,你可以像查询普通表一样进行JOIN、聚合等操作。
假设我们的电商数据库 ecommerce 中有两张核心表:orders(订单表)和 order_details(订单明细表)。我们的目标是将它们的变化实时同步到Kafka,并关联成一张宽表。
首先,在Flink SQL Client中创建CDC源表:
-- 创建订单表CDC源
CREATE TABLE mysql_orders (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status INT,
create_time TIMESTAMP(3),
update_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql',
'port' = '3306',
'username' = 'flinkuser',
'password' = 'flinkpw',
'database-name' = 'ecommerce',
'table-name' = 'orders',
'server-time-zone' = 'Asia/Shanghai',
'debezium.snapshot.mode' = 'initial' -- 首次启动时做全量快照
);
-- 创建订单明细表CDC源
CREATE TABLE mysql_order_details (
detail_id BIGINT,
order_id BIGINT,
product_id BIGINT,
product_name STRING,
price DECIMAL(10, 2),
quantity INT,
PRIMARY KEY (detail_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql',
'port' = '3306',
'username' = 'flinkuser',
'password' = 'flinkpw',
'database-name' = 'ecommerce',
'table-name' = 'order_details'
);
接下来,我们将这两张流表进行关联,并写入Kafka,形成订单宽表流:
-- 创建Kafka目标表
CREATE TABLE kafka_order_wide (
order_id BIGINT,
user_id BIGINT,
product_id BIGINT,
product_name STRING,
quantity INT,
price DECIMAL(10, 2),
total_amount DECIMAL(10, 2),
order_status INT,
event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'order_wide',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json',
'json.ignore-parse-errors' = 'true'
);
-- 将关联后的宽表数据插入Kafka
INSERT INTO kafka_order_wide
SELECT
o.order_id,
o.user_id,
d.product_id,
d.product_name,
d.quantity,
d.price,
o.total_amount,
o.order_status,
GREATEST(o.update_time, o.create_time) as event_time -- 取最近时间作为事件时间
FROM mysql_orders o
JOIN mysql_order_details d
ON o.order_id = d.order_id;
执行这段SQL后,一个实时的订单宽表数据流就会持续写入Kafka的 order_wide 主题。这里有几个关键点需要注意:
提示:Flink CDC在
initial模式下,会先读取全量数据快照,然后无缝切换到binlog增量读取。这保证了数据不丢失。对于大表,可以考虑使用initial或schema_only模式以避免长时间锁表。
- 关联操作:在流上进行JOIN,Flink内部会维护状态。需要根据业务逻辑合理设置状态的TTL(生存时间),防止状态无限膨胀。
- 事件时间与水位线:我们定义了
WATERMARK,这是后续进行基于时间的窗口聚合的关键。它告诉系统事件时间的进展,允许处理一定程度的数据乱序。 - Exactly-Once保证:要确保从CDC到Kafka的端到端一致性,需要开启Flink的Checkpoint机制,并配置Kafka生产者的事务属性。
4. 实时ETL与聚合:流式宽表与DWD层建设
数据进入Kafka后,原始的订单宽表流(可以视为实时ODS层)还需要进一步清洗、维度关联和轻度汇总,形成实时DWD(明细数据层)和DWS(汇总数据层)。这是实时数仓业务逻辑最密集的部分。
场景一:实时过滤与清洗 假设我们只关心状态为“已支付”(order_status=2)的订单,并且需要过滤掉金额异常(如total_amount <= 0)的数据。
CREATE TABLE kafka_order_wide_cleaned (
-- ... 字段同前 ...
) WITH (
'connector' = 'kafka',
'topic' = 'order_wide_cleaned',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json'
);
INSERT INTO kafka_order_wide_cleaned
SELECT * FROM kafka_order_wide
WHERE order_status = 2 AND total_amount > 0;
场景二:实时关联维度信息(如用户画像) 维度数据通常存储在HBase、Redis或维度表中。这里演示如何通过Flink SQL的LOOKUP JOIN关联静态的用户维度表(假设已预加载到Flink状态中)。
-- 假设有一张静态的用户维度表(可从MySQL周期性加载)
CREATE TABLE dim_user (
user_id BIGINT,
user_level INT,
city STRING,
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://mysql:3306/ecommerce',
'table-name' = 'dim_user',
'username' = 'flinkuser',
'password' = 'flinkpw',
'lookup.cache.max-rows' = '10000', -- 缓存最大行数
'lookup.cache.ttl' = '10 min' -- 缓存过期时间
);
-- 实时关联用户维度
CREATE TABLE kafka_order_wide_with_user (
order_id BIGINT,
user_id BIGINT,
user_level INT,
city STRING,
product_name STRING,
total_amount DECIMAL(10, 2),
event_time TIMESTAMP(3)
) WITH ( ... );
INSERT INTO kafka_order_wide_with_user
SELECT
o.order_id,
o.user_id,
u.user_level,
u.city,
o.product_name,
o.total_amount,
o.event_time
FROM kafka_order_wide_cleaned o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.event_time AS u
ON o.user_id = u.user_id;
场景三:核心实时聚合指标计算 这是实时数仓的价值所在。我们计算每分钟每个城市的销售总额和订单数。
CREATE TABLE iceberg_dws_city_sales_minute ( -- 目标表改为Iceberg
city STRING,
window_start TIMESTAMP(3),
window_end TIMESTAMP(3),
sales_amount DECIMAL(10, 2),
order_count BIGINT,
PRIMARY KEY (city, window_start) NOT ENFORCED
) WITH (
'connector' = 'iceberg',
'catalog-name' = 'hive_prod',
'catalog-type' = 'hive',
'warehouse' = 'hdfs://localhost:9000/user/iceberg/warehouse',
'format-version' = '2'
);
-- 使用窗口TVF(Table-Valued Function)语法,更灵活强大
INSERT INTO iceberg_dws_city_sales_minute
SELECT
city,
window_start,
window_end,
SUM(total_amount) as sales_amount,
COUNT(DISTINCT order_id) as order_count
FROM TABLE(
TUMBLE(TABLE kafka_order_wide_with_user, DESCRIPTOR(event_time), INTERVAL '1' MINUTE)
)
GROUP BY city, window_start, window_end;
这段代码将聚合结果直接写入了Apache Iceberg表。这是3.0架构的关键一步:实时聚合结果不再只存在于Kafka主题或内存状态中,而是以结构化表格的形式持久化到分布式文件系统(如HDFS或S3)。这样做的好处是:
- 数据可回溯:任何时间点的聚合结果都被保存下来,可以随时查询历史任意分钟的数据。
- 支持批处理:如果需要重新计算历史某天的数据(比如维度表变更),可以直接以Iceberg表为源启动一个批处理Flink作业,而无需从Kafka重放。
- 统一服务层:下游的OLAP引擎(如Doris)可以直接读取Iceberg表进行加速查询。
5. 数据服务与可视化:基于实时OLAP的即席查询
实时数据经过ETL和聚合,写入Iceberg后,最后一步是如何高效地提供给业务方使用。直接查询Iceberg表虽然可行,但其延迟通常在秒到十秒级,对于交互式仪表盘来说还不够快。因此,我们需要引入专用的实时OLAP引擎作为加速层。
这里以Apache Doris为例,展示如何构建数据服务层。
第一步:在Doris中创建外部表映射Iceberg表 Doris支持通过Multi-Catalog功能直接映射外部数据源,无需数据导入。
-- 在Doris中执行
CREATE CATALOG iceberg_catalog PROPERTIES (
"type"="iceberg",
"iceberg.catalog.type"="hive",
"hive.metastore.uris"="thrift://hive-metastore:9083"
);
-- 切换catalog
SWITCH iceberg_catalog;
-- 此时可以直接查询Iceberg表
SELECT * FROM `ecommerce`.`iceberg_dws_city_sales_minute` WHERE city='北京' ORDER BY window_start DESC LIMIT 10;
第二步:创建Doris物化视图(Rollup)进行实时预聚合 对于最热门的查询,如“今日实时销售大盘”,我们可以在Doris内部创建物化视图,实现亚秒级响应。
-- 在Doris中为明细表创建物化视图(需先将数据导入Doris内部表)
-- 假设我们已将订单宽表通过Flink-Doris-Connector实时导入到Doris表`doris_order_wide`
CREATE MATERIALIZED VIEW dws_realtime_dashboard
BUILD IMMEDIATE -- 立即构建
REFRESH ASYNC -- 异步刷新,与底层数据变更保持同步
KEY (product_id, hour)
DISTRIBUTED BY HASH(product_id)
AS
SELECT
product_id,
DATE_TRUNC('hour', event_time) as hour,
SUM(total_amount) as hourly_sales,
COUNT(DISTINCT order_id) as hourly_orders
FROM doris_order_wide
WHERE event_time >= TODAY()
GROUP BY product_id, DATE_TRUNC('hour', event_time);
第三步:配置BI工具连接 最后,在Superset、DataEase等BI工具中,添加Doris数据源,就可以轻松地拖拽生成实时数据看板了。你可以创建诸如“实时销售趋势图”、“地域销售热力图”、“畅销商品排行榜”等可视化组件。
注意:在生产环境中,需要密切关注Doris集群的负载和查询性能。对于复杂的Ad-hoc查询,建议通过设置合适的分区、分桶策略,以及创建针对性的物化视图来优化。同时,要建立监控告警,关注数据从Iceberg到Doris的同步延迟。
走到这里,一个完整的电商实时数仓3.0系统已经搭建完毕。从MySQL的变更数据捕获,到Flink流式处理与聚合,再到Iceberg统一存储,最后通过Doris提供高速查询服务,形成了一条高效、稳定、可扩展的数据流水线。这套架构的优势在于,它用一套技术栈解决了实时和离线两种场景,极大地降低了开发和运维的复杂度。当然,每个环节都有大量的调优细节,比如Flink状态管理、Iceberg小文件合并、Doris查询优化等,这些都需要在具体实践中不断摸索和积累。我在最初上线类似系统时,就曾因为没设置好状态TTL导致作业崩溃,也遇到过Iceberg分区策略不当引发的查询慢问题。我的经验是,先让管道跑通,再针对性能瓶颈逐个击破,同时建立完善的监控体系,这样才能让实时数仓真正稳定、可靠地驱动业务决策。
更多推荐
所有评论(0)