【实时数据仓库建模】Kafka+Flink实战
【实时数据仓库建模】Kafka+Flink实战
一、引言
-
1、核心:通过
Kafka+Flink构建实时数据仓库,解决传统数据仓库“慢”的问题,将数据从产生到分析的延迟缩短到秒级或分钟级 -
Kafka:实时数据管道,负责收集、存储和传输海量流数据,具备高吞吐量、低延迟、可持久化等特性
-
Flink:实时计算引擎,负责对Kafka中的流数据进行清洗、转换、聚合,支持 精确一次(Exactly-Once) 语义,确保数据一致性。
-
2、最终效果展示
本文参考链接:https://blog.csdn.net/2501_91888447/article/details/155677083
本文与参考链接的区别是:将参考链接中每一步进行实现实战并截图说明,其中python构造数据部分换成Nodered实现
本文将以电商实时销售额监控为例,构建一个端到端的实时数据仓库。最终实现:
- 实时 dashboard:展示每小时销售额、订单量、TOP 10商品(更新频率:1分钟)
- 数据链路:订单数据从产生(Kafka生产者)→ 存储(Kafka topic)→ 处理(Flink)→ 分析(ClickHouse)→ 展示(Grafana),全程延迟≤1分钟。
- 3、环境与工具清单

- 4、 环境搭建步骤(简要)
- 安装Kafka与ZooKeeper:参考Kafka官方文档;
- 安装ClickHouse:参考ClickHouse官方文档;
- 安装Grafana:参考Grafana官方文档
- 安装MySQL:参考MySQL官方文档
- 安装Flink:参考Flink官方文档;
注意:
这里的Flink需要开启SQL Gateway服务,并支持Kafka和ClickHouse插件
参考实战链接:【Flink on Kubernetes部署详细教程】
- 安装Nodered:参考Nodered官方文档
二、核心步骤
2.1 需求分析与建模
- 需求分析:本次实战的需求是电商平台实时销售额监控,具体如下:
实时统计:每小时的总销售额、总订单量;
TOP商品:每小时销售额前10的商品;
数据延迟:从订单产生到 dashboard 展示≤1分钟;
数据准确性:确保数据不丢失、不重复(精确一次语义)。 - 数据建模:实时数据仓库分层设计,强调低延迟和流处理特性。本文采用经典的“三层模型”
| 层级 | 全称 | 作用说明 | 存储介质 |
|---|---|---|---|
| ODS | 操作数据存储 | 保留原始数据(如订单、用户行为),不做任何清洗,方便回溯 | Kafka Topic |
| DWD | 数据仓库明细层 | 对ODS层数据进行清洗(过滤无效数据)、补全(关联维度表),生成干净的明细数据 | Kafka Topic |
| DWS | 数据仓库汇总层 | 按主题(如"销售额"、“订单量”)进行实时聚合,生成汇总数据(支持快速查询) | ClickHouse 表 |
- ODS层:原始数据存储
ODS层的目标是保存原始数据,避免数据丢失。本次实战中,ODS层对应Kafka的order_topic,存储从生产者发送的原始订单数据。

- DWD层:明细数据清洗
DWD层的目标是干净、完整,对应Kafka的dwd_order_topic,存储清洗后的明细数据。主要做两件事:
过滤无效数据:比如取消的订单(status=‘canceled’)不需要统计;
关联维度表:补全商品名称(product_name),方便后续分析。 - DWS层:汇总数据计算
DWS层的目标是快速查询,按时间窗口(如1小时)对DWD层数据进行聚合,生成汇总结果。本次实战中,DWS层包含两张表:
dws_hourly_sales:每小时总销售额、总订单量;
dws_hourly_top_products:每小时销售额前10的商品。
这些汇总数据存储在ClickHouse中,因为ClickHouse是列式存储数据库,支持秒级查询,非常适合实时分析。
2.2 数据表创建
- 创建MySQL维度表:
CREATE DATABASE db_test;
CREATE TABLE db_test.product (
product_id INT PRIMARY KEY AUTO_INCREMENT,
product_name VARCHAR(100) NOT NULL,
price DECIMAL(10,2) NOT NULL
);
– 插入测试数据
INSERT INTO db_test.product (product_name, price) VALUES
(‘iPhone 15’, 7999.00),
(‘MacBook Pro’, 14999.00),
(‘iPad Pro’, 6999.00),
(‘Apple Watch Series 9’, 2999.00),
(‘AirPods Pro 2’, 1999.00);
2.3 Kafka生产者 -> Nodered模拟订单数据
首先,我们需要一个Kafka生产者,模拟电商平台的订单生成。这里用Nodered实现:


此时能看到Kafka的order_topic中已经有了模拟的订单数据:

点击跳转查看源码
回到目录
2.4 Flink处理 -> 从ODS到DWD再到DWS
Flink是本次实战的“核心引擎”,负责处理从Kafka读取的数据,完成DWD层清洗和DWS层聚合。我们使用Flink SQL(更易上手,适合批量处理)来实现。
- 在DBEaver连接Flink SQL客户端(或直接在Flink SQL客户端执行sql)中创建ODS层表ods_order,关联Kafka的order_topic:
CREATE TABLE IF NOT EXISTS ods_order (
order_id INT,
user_id INT,
product_id INT,
amount DECIMAL(10, 2),
create_time STRING,
status STRING,
proc_time AS PROCTIME(),
event_time AS TO_TIMESTAMP(create_time, 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z'''),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'order_topic',
'properties.bootstrap.servers' = '10.30.2.95:9092',
'properties.group.id' = 'flink_ods_order_group',
'format' = 'json',
'scan.startup.mode' = 'latest-offset'
);
order_id:订单ID user_id: 用户ID product_id:商品ID amount :订单金额 create_time :订单创建时间status :订单状态
event_time:事件时间,即订单的实际创建时间(从create_time字段转换而来)WATERMARK:水印,用于处理事件时间的延迟(比如订单数据因网络问题延迟5秒到达,水印会等待5秒再关闭窗口,确保数据不丢失)。
- MySQL商品信息维度表(与MySQL结构一致):
CREATE TABLE IF NOT EXISTS product_dim (
product_id INT PRIMARY KEY NOT ENFORCED,
product_name STRING,
price DECIMAL(10, 2)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://10.30.2.85:3306/db_test',
'table-name' = 'product',
'username' = 'jxdev',
'password' = 'ahjuxin_2026' ,
'lookup.cache.max-rows' = '1000',
'lookup.cache.ttl' = '1h'
);
- DWD层表dwd_order:清洗后的明细数据并写入kafka
CREATE TABLE IF NOT EXISTS dwd_order (
order_id INT,
user_id INT,
product_id INT,
product_name STRING,
amount DECIMAL(10, 2),
create_time STRING,
status STRING,
event_time TIMESTAMP(3),
-- 关键:定义事件时间与水印,把event_time转为时间属性
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'dwd_order_topic',
'properties.bootstrap.servers' = '10.30.2.95:9092',
'properties.group.id' = 'flink_dws_hourly_sales_group',
'format' = 'json',
'sink.partitioner' = 'round-robin'
);
-- 插入数据到DWD层(过滤取消订单+关联商品信息维度表)
INSERT INTO dwd_order
SELECT o.order_id,o.user_id,o.product_id,p.product_name, o.amount,o.create_time,o.status, o.event_time
FROM ods_order o
JOIN product_dim FOR SYSTEM_TIME AS OF o.proc_time p ON o.product_id = p.product_id
WHERE o.status <> 'canceled';
FOR SYSTEM_TIME AS OF o.proc_time:Flink的时态表关联(Temporal Table Join),用于关联维度表的快照(避免维度数据更新导致的不一致);
dwd_order_topic:DWD层数据存储的Kafka Topic,后续可以用于其他分析(如用户行为分析)。
- 创建DWS层表:实时聚合计算
DWS层是实时数据仓库的“出口”,负责生成汇总数据,供可视化工具查询。本次实战中,创建两张DWS层表:
① 每小时销售额与订单量(dws_hourly_sales)
② 每小时TOP 10商品(dws_hourly_top_products)
-- 5、创建DWS层表:存储每小时销售额汇总数据,并写入ClickHouse表
CREATE TABLE IF NOT EXISTS dws_hourly_sales (
`hour` STRING, -- 小时(格式:yyyy-MM-dd HH:00)
total_sales DECIMAL(10, 2), -- 总销售额
total_orders INT, -- 总订单量
window_start TIMESTAMP(3), -- 窗口开始时间
window_end TIMESTAMP(3) -- 窗口结束时间
) WITH (
'connector' = 'clickhouse',
'url' = 'jdbc:ch://10.30.2.98:8123',
'database-name' = 'bench',
'table-name' = 'dws_hourly_sales',
'username' = 'default',
'password' = 'Huayu_2025',
'sink.batch-size' = '1000',
'sink.flush-interval' = '1000'
);
-- 插入数据到DWS层(每小时汇总销售额与订单量)
INSERT INTO dws_hourly_sales
SELECT
DATE_FORMAT(window_start, 'yyyy-MM-dd HH:00') AS `hour`,
SUM(amount) AS total_sales,
CAST(COUNT(DISTINCT order_id) AS INT) AS total_orders,
window_start,
window_end
FROM TABLE( TUMBLE(TABLE dwd_order, DESCRIPTOR(event_time), INTERVAL '1' HOUR) )
GROUP BY window_start, window_end;
-- 6、创建DWS层表:存储每小时TOP 10商品数据,写入ClickHouse
CREATE TABLE dws_hourly_top_products (
`hour` STRING, -- 小时(格式:yyyy-MM-dd HH:00)
product_id INT, -- 商品ID
product_name STRING, -- 商品名称
sales DECIMAL(10, 2), -- 商品销售额
`rank` INT, -- 排名(1-10)
window_start TIMESTAMP(3), -- 窗口开始时间
window_end TIMESTAMP(3), -- 窗口结束时间
PRIMARY KEY (`hour`, product_id) NOT ENFORCED
) WITH (
'connector' = 'clickhouse',
'url' = 'jdbc:ch://10.30.2.98:8123',
'database-name' = 'bench',
'table-name' = 'dws_hourly_top_products',
'username' = 'default',
'password' = 'Huayu_2025',
'sink.batch-size' = '1000',
'sink.flush-interval' = '1000'
);
-- 插入数据到DWS层(每小时TOP 10商品)
INSERT INTO dws_hourly_top_products
SELECT
DATE_FORMAT(window_start, 'yyyy-MM-dd HH:00') AS `hour`,
product_id,
product_name,
sales,
CAST(rn AS INT) as `rank`,
window_start,
window_end
FROM (
SELECT
product_id,
product_name,
SUM(amount) AS sales,
window_start,
window_end,
ROW_NUMBER() OVER (PARTITION BY window_start ORDER BY SUM(amount) DESC) AS rn
FROM TABLE(
TUMBLE(TABLE dwd_order, DESCRIPTOR(event_time), INTERVAL '1' HOUR)
)
GROUP BY product_id, product_name, window_start, window_end
)
WHERE rn <= 10;
说明:
- 滚动窗口(Tumbling Window):每1小时一个窗口,窗口之间不重叠(如14:00-15:00、15:00-16:00),适合固定时间间隔的汇总;
- RANK()函数:用于计算每个商品在窗口内的排名(降序),PARTITION BY window_start表示按窗口分组,ORDER BY SUM(amount) DESC表示按销售额降序排列;
- ClickHouse存储:ClickHouse的列式存储和向量查询引擎,使得汇总数据的查询速度非常快(秒级),适合实时展示。
- 验证表是否创建成功:
show tables;

执行成功后,可以看到flink客户端有正在运行中的任务:

三、源码
3.1 Nodered模拟订单发送到Kafka源码
[
{
"id": "67428b8f6878b23b",
"type": "inject",
"z": "067b766577a0b317",
"name": "",
"props": [
{
"p": "payload"
},
{
"p": "topic",
"vt": "str"
}
],
"repeat": "",
"crontab": "",
"once": false,
"onceDelay": 0.1,
"topic": "",
"payload": "",
"payloadType": "date",
"x": 90,
"y": 60,
"wires": [
[
"51cdba257295fe80"
]
]
},
{
"id": "51cdba257295fe80",
"type": "function",
"z": "067b766577a0b317",
"name": "初始化",
"func": "// 商品列表(与MySQL维度表一致)\nmsg.products = [\n {'product_id': 1, 'product_name': 'iPhone 15', 'price': 7999.00},\n {'product_id': 2, 'product_name': 'MacBook Pro', 'price': 14999.00},\n {'product_id': 3, 'product_name': 'iPad Pro', 'price': 6999.00},\n {'product_id': 4, 'product_name': 'Apple Watch Series 9', 'price': 2999.00},\n {'product_id': 5, 'product_name': 'AirPods Pro 2', 'price': 1999.00}\n]\nmsg.orderStatusList = ['completed', 'pending', 'canceled'];//订单状态\nreturn msg;",
"outputs": 1,
"timeout": 0,
"noerr": 0,
"initialize": "",
"finalize": "",
"libs": [],
"x": 230,
"y": 60,
"wires": [
[
"fc8302a4d3979c14"
]
]
},
{
"id": "8582719c984a6e2b",
"type": "function",
"z": "067b766577a0b317",
"name": "生成模拟订单",
"func": "// 生成模拟订单\nconst product = msg.products[randomInt(0,4)];//订单选购商品\nmsg.payload = {\n 'order_id': randomInt(100000, 999999), // 随机订单ID\n 'user_id': randomInt(1, 10000), // 随机用户ID\n 'product_id': product.product_id,// 商品ID(来自商品列表)\n 'amount': Number((randomInt(1, 5) * product.price).toFixed(2)), // 订单金额(数量×单价)\n 'create_time': new Date(), // 订单创建时间(当前时间)\n 'status': msg.orderStatusList[randomInt(0, 2)] // 随机订单状态\n}\nreturn msg;\n/** 生成区间[min,max]内的随机数int */\nfunction randomInt(min, max) {\n const ran = Math.floor(Math.random() * (max - min + 1)) + min;\n if (ran >= max){ return max; }\n if (ran <= min) { return min; }\n return ran;\n}",
"outputs": 1,
"timeout": 0,
"noerr": 0,
"initialize": "",
"finalize": "",
"libs": [],
"x": 380,
"y": 140,
"wires": [
[
"52d6654949bb37c1",
"8fc61fbcdcc14c6b"
]
]
},
{
"id": "8fc61fbcdcc14c6b",
"type": "delay",
"z": "067b766577a0b317",
"name": "随机延迟1-5秒",
"pauseType": "random",
"timeout": "5",
"timeoutUnits": "seconds",
"rate": "1",
"nbRateUnits": "1",
"rateUnits": "second",
"randomFirst": "1",
"randomLast": "5",
"randomUnits": "seconds",
"drop": false,
"allowrate": false,
"outputs": 1,
"x": 620,
"y": 140,
"wires": [
[
"fc8302a4d3979c14"
]
]
},
{
"id": "fc8302a4d3979c14",
"type": "counter-loop",
"z": "067b766577a0b317",
"name": "",
"counter": "cid",
"counterType": "msg",
"reset": false,
"resetValue": "value-null",
"initial": 0,
"initialType": "num",
"operator": "lt",
"termination": "100",
"terminationType": "num",
"increment": 1,
"incrementType": "num",
"x": 440,
"y": 60,
"wires": [
[
"8ed5ca6d35c83157"
],
[
"8582719c984a6e2b"
]
]
},
{
"id": "52d6654949bb37c1",
"type": "rdkafka out",
"z": "067b766577a0b317",
"name": "",
"topic": "order_topic",
"key": "",
"partition": -1,
"broker": "c331c6a48d89a75a",
"x": 610,
"y": 180,
"wires": []
},
{
"id": "8ed5ca6d35c83157",
"type": "function",
"z": "067b766577a0b317",
"name": "结束",
"func": "node.warn(\"模拟订单发送结束,共产生\"+msg.cid+\"条订单\");\nreturn msg;",
"outputs": 1,
"timeout": 0,
"noerr": 0,
"initialize": "",
"finalize": "",
"libs": [],
"x": 590,
"y": 40,
"wires": [
[]
]
},
{
"id": "c331c6a48d89a75a",
"type": "kafka-broker",
"broker": "10.30.2.95:9092",
"clientid": ""
}
]
更多推荐
所有评论(0)