一、引言

  • 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、环境与工具清单
    1
  • 4、 环境搭建步骤(简要)
  • 安装Kafka与ZooKeeper:参考Kafka官方文档;
  • 安装ClickHouse:参考ClickHouse官方文档;
  • 安装Grafana:参考Grafana官方文档
  • 安装MySQL:参考MySQL官方文档
  • 安装Flink:参考Flink官方文档;

注意:
  这里的Flink需要开启SQL Gateway 服务,并支持Kafka和ClickHouse插件
  参考实战链接:【Flink on Kubernetes部署详细教程

  • 安装Nodered:参考Nodered官方文档

回到目录

二、核心步骤

2.1 需求分析与建模

  1. 需求分析:本次实战的需求是电商平台实时销售额监控,具体如下:
      实时统计:每小时的总销售额、总订单量;
      TOP商品:每小时销售额前10的商品;
      数据延迟:从订单产生到 dashboard 展示≤1分钟;
      数据准确性:确保数据不丢失、不重复(精确一次语义)。
  2. 数据建模:实时数据仓库分层设计,强调低延迟和流处理特性。本文采用经典的“三层模型”
层级全称作用说明存储介质
ODS操作数据存储保留原始数据(如订单、用户行为),不做任何清洗,方便回溯Kafka Topic
DWD数据仓库明细层对ODS层数据进行清洗(过滤无效数据)、补全(关联维度表),生成干净的明细数据Kafka Topic
DWS数据仓库汇总层按主题(如"销售额"、“订单量”)进行实时聚合,生成汇总数据(支持快速查询)ClickHouse 表
  • ODS层:原始数据存储
    ODS层的目标是保存原始数据,避免数据丢失。本次实战中,ODS层对应Kafka的order_topic,存储从生产者发送的原始订单数据。
    2
  • 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 数据表创建

  1. 创建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实现:
1
2
此时能看到Kafka的order_topic中已经有了模拟的订单数据:
3
点击跳转查看源码
回到目录

2.4 Flink处理 -> 从ODS到DWD再到DWS

  Flink是本次实战的“核心引擎”,负责处理从Kafka读取的数据,完成DWD层清洗和DWS层聚合。我们使用Flink SQL(更易上手,适合批量处理)来实现。

  1. 在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秒再关闭窗口,确保数据不丢失)。
  1. 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'
);
  1. 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,后续可以用于其他分析(如用户行为分析)。

  1. 创建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的列式存储和向量查询引擎,使得汇总数据的查询速度非常快(秒级),适合实时展示。
  1. 验证表是否创建成功:show tables;
    4
    执行成功后,可以看到flink客户端有正在运行中的任务:
    5

回到目录

三、源码

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": ""
    }
]

回到目录

更多推荐