实时数据仓库建模:Kafka+Flink实战案例

引言

痛点引入:传统数据仓库的“慢”问题

在数字化时代,企业对数据的需求早已从“事后分析”转向“实时决策”。比如:

  • 电商平台需要实时监控每小时销售额,及时调整促销策略;
  • 物流企业需要实时追踪订单状态,避免延误;
  • 金融机构需要实时检测交易欺诈,降低风险。

但传统数据仓库(如Hive+Spark离线数仓)的“T+1”或“小时级”延迟,根本无法满足这些需求。比如,当运营人员看到昨天的销售额报表时,可能已经错过了最佳的调整时机。

解决方案:Kafka+Flink构建实时数据仓库

实时数据仓库的核心目标是将数据从产生到分析的延迟缩短到秒级或分钟级。而Kafka+Flink的组合,正是实现这一目标的“黄金搭档”:

  • Kafka:作为实时数据管道,负责收集、存储和传输海量流数据(如订单、用户行为),具备高吞吐量、低延迟、可持久化等特性;
  • Flink:作为实时计算引擎,负责对Kafka中的流数据进行清洗、转换、聚合,支持**精确一次(Exactly-Once)**语义,确保数据一致性。

最终效果展示

本文将以电商实时销售额监控为例,构建一个端到端的实时数据仓库。最终实现:

  • 实时 dashboard:展示每小时销售额、订单量、TOP 10商品(更新频率:1分钟);
  • 数据链路:订单数据从产生(Kafka生产者)→ 存储(Kafka topic)→ 处理(Flink)→ 分析(ClickHouse)→ 展示(Grafana),全程延迟≤1分钟。

准备工作

1. 环境与工具清单

工具/组件 版本 作用说明
Kafka 2.8.0+ 实时数据管道,存储原始订单数据和处理后的明细数据
Flink 1.17.0+ 实时计算引擎,负责数据清洗、聚合
ZooKeeper 3.6.3+ Kafka依赖的分布式协调服务
ClickHouse 23.8.0+ 列式存储数据库,用于存储汇总数据(支持快速查询)
Grafana 10.0.0+ 可视化工具,展示实时 dashboard
MySQL 8.0+ 存储维度数据(如商品信息)
Python 3.8+ 模拟订单数据生产者

2. 基础知识要求

  • 了解Kafka的核心概念(Topic、Partition、Producer/Consumer);
  • 了解Flink的基本用法(DataStream API、Table API/SQL);
  • 了解实时数据仓库的分层模型(ODS、DWD、DWS,下文会详细解释)。

3. 环境搭建步骤(简要)

  1. 安装Kafka与ZooKeeper:参考Kafka官方文档
  2. 安装Flink:参考Flink官方文档
  3. 安装ClickHouse:参考ClickHouse官方文档
  4. 安装Grafana:参考Grafana官方文档
  5. 创建MySQL维度表
    CREATE DATABASE shop;
    USE shop;
    CREATE TABLE product (
        product_id INT PRIMARY KEY AUTO_INCREMENT,
        product_name VARCHAR(100) NOT NULL,
        price DECIMAL(10,2) NOT NULL
    );
    -- 插入测试数据
    INSERT INTO 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);
    

核心步骤:从需求到实现

一、需求分析:明确要解决的问题

本次实战的需求是电商平台实时销售额监控,具体要求:

  1. 实时统计:每小时的总销售额、总订单量;
  2. TOP商品:每小时销售额前10的商品;
  3. 数据延迟:从订单产生到 dashboard 展示≤1分钟;
  4. 数据准确性:确保数据不丢失、不重复(精确一次语义)。

二、数据建模:实时数据仓库分层设计

实时数据仓库的分层模型与离线数仓类似,但更强调低延迟流处理特性。本文采用经典的“三层模型”:

层级 全称 作用说明 存储介质
ODS 操作数据存储 保留原始数据(如订单、用户行为),不做任何清洗,方便回溯 Kafka Topic
DWD 数据仓库明细层 对ODS层数据进行清洗(过滤无效数据)、补全(关联维度表),生成干净的明细数据 Kafka Topic
DWS 数据仓库汇总层 按主题(如“销售额”“订单量”)进行实时聚合,生成汇总数据(支持快速查询) ClickHouse
1. ODS层:原始数据存储

ODS层的目标是**“原汁原味”保存原始数据**,避免数据丢失。本次实战中,ODS层对应Kafka的order_topic,存储从生产者发送的原始订单数据。

订单数据字段说明

字段名 类型 说明
order_id INT 订单ID(唯一标识)
user_id INT 用户ID
product_id INT 商品ID
amount DECIMAL(10,2) 订单金额(数量×单价)
create_time TIMESTAMP 订单创建时间
status STRING 订单状态(completed:完成;pending:待支付;canceled:取消)
2. DWD层:明细数据清洗

DWD层的目标是**“干净、完整”**,主要做两件事:

  • 过滤无效数据:比如取消的订单(status=‘canceled’)不需要统计;
  • 关联维度表:补全商品名称(product_name),方便后续分析。

DWD层对应Kafka的dwd_order_topic,存储清洗后的明细数据。

3. DWS层:汇总数据计算

DWS层的目标是**“快速查询”**,按时间窗口(如1小时)对DWD层数据进行聚合,生成汇总结果。本次实战中,DWS层包含两张表:

  • dws_hourly_sales:每小时总销售额、总订单量;
  • dws_hourly_top_products:每小时销售额前10的商品。

这些汇总数据存储在ClickHouse中,因为ClickHouse是列式存储数据库,支持秒级查询,非常适合实时分析。

三、数据Pipeline实现:从Kafka到Flink再到ClickHouse

1. 步骤1:Kafka生产者——模拟订单数据

首先,我们需要一个Kafka生产者,模拟电商平台的订单生成。这里用Python的kafka-python库实现:

代码示例(producer.py)

from kafka import KafkaProducer
import json
import time
import random
from datetime import datetime

# Kafka配置
bootstrap_servers = 'localhost:9092'
topic = 'order_topic'

# 商品列表(与MySQL维度表一致)
products = [
    {'product_id': 1, 'product_name': 'iPhone 15', 'price': 7999.00},
    {'product_id': 2, 'product_name': 'MacBook Pro', 'price': 14999.00},
    {'product_id': 3, 'product_name': 'iPad Pro', 'price': 6999.00},
    {'product_id': 4, 'product_name': 'Apple Watch Series 9', 'price': 2999.00},
    {'product_id': 5, 'product_name': 'AirPods Pro 2', 'price': 1999.00}
]

# 初始化Kafka生产者
producer = KafkaProducer(
    bootstrap_servers=bootstrap_servers,
    value_serializer=lambda v: json.dumps(v).encode('utf-8')  # 将数据序列化为JSON格式
)

def generate_order():
    """生成模拟订单数据"""
    selected_product = random.choice(products)
    return {
        'order_id': random.randint(100000, 999999),  # 随机订单ID
        'user_id': random.randint(1, 10000),         # 随机用户ID
        'product_id': selected_product['product_id'],# 商品ID(来自商品列表)
        'amount': round(random.randint(1, 5) * selected_product['price'], 2),  # 订单金额(数量×单价)
        'create_time': datetime.now().strftime('%Y-%m-%d %H:%M:%S'),  # 订单创建时间(当前时间)
        'status': random.choice(['completed', 'pending', 'canceled'])  # 随机订单状态
    }

if __name__ == '__main__':
    try:
        while True:
            order = generate_order()
            producer.send(topic, value=order)  # 发送订单到Kafka
            print(f"发送订单成功:{order}")
            time.sleep(random.randint(1, 3))  # 模拟订单产生的随机延迟(1-3秒)
    except KeyboardInterrupt:
        print("停止发送订单")
    finally:
        producer.close()  # 关闭生产者

运行生产者

pip install kafka-python
python producer.py

此时,Kafka的order_topic中已经有了模拟的订单数据,可以用Kafka自带的消费者工具验证:

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic order_topic --from-beginning
2. 步骤2:Flink处理——从ODS到DWD再到DWS

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

(1)创建ODS层表:读取Kafka原始数据

首先,在Flink SQL客户端中创建ODS层表ods_order,关联Kafka的order_topic

Flink SQL代码

-- 创建ODS层表(关联Kafka原始订单数据)
CREATE TABLE ods_order (
    order_id INT,                  -- 订单ID
    user_id INT,                   -- 用户ID
    product_id INT,                -- 商品ID
    amount DECIMAL(10, 2),         -- 订单金额
    create_time STRING,            -- 订单创建时间(字符串类型,后续转换为TIMESTAMP)
    status STRING,                 -- 订单状态
    -- 事件时间(用于窗口计算):从create_time字段提取,转换为TIMESTAMP类型
    event_time AS TO_TIMESTAMP(create_time, 'yyyy-MM-dd HH:mm:ss'),
    -- 水印(Watermark):处理事件时间的延迟,允许5秒延迟
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',                          -- 连接器类型:Kafka
    'topic' = 'order_topic',                        -- 关联的Kafka Topic
    'properties.bootstrap.servers' = 'localhost:9092',  -- Kafka集群地址
    'properties.group.id' = 'flink_ods_order_group',    -- 消费者组ID
    'format' = 'json',                              -- 数据格式:JSON
    'scan.startup.mode' = 'latest-offset'           -- 消费起始位置:最新偏移量(避免重复消费历史数据)
);

说明

  • event_time:事件时间,即订单的实际创建时间(从create_time字段转换而来);
  • WATERMARK:水印,用于处理事件时间的延迟(比如订单数据因网络问题延迟5秒到达,水印会等待5秒再关闭窗口,确保数据不丢失)。
(2)创建DWD层表:清洗与维表关联

接下来,创建DWD层表dwd_order,完成两件事:

  • 过滤取消的订单(status != 'canceled');
  • 关联MySQL维度表product,补全商品名称(product_name)。

首先,创建MySQL维度表的连接:

-- 创建MySQL维度表(商品信息)
CREATE TABLE product_dim (
    product_id INT PRIMARY KEY,    -- 商品ID(主键)
    product_name STRING,           -- 商品名称
    price DECIMAL(10, 2)           -- 商品单价
) WITH (
    'connector' = 'jdbc',                          -- 连接器类型:JDBC
    'url' = 'jdbc:mysql://localhost:3306/shop',    -- MySQL地址(数据库:shop)
    'table-name' = 'product',                       -- 关联的MySQL表名
    'username' = 'root',                            -- MySQL用户名
    'password' = 'root'                             -- MySQL密码
);

然后,创建DWD层表并插入数据:

-- 创建DWD层表(存储清洗后的明细数据)
CREATE TABLE dwd_order (
    order_id INT,                  -- 订单ID
    user_id INT,                   -- 用户ID
    product_id INT,                -- 商品ID
    product_name STRING,           -- 商品名称(来自维度表)
    amount DECIMAL(10, 2),         -- 订单金额
    create_time STRING,            -- 订单创建时间
    status STRING,                 -- 订单状态
    event_time TIMESTAMP(3)        -- 事件时间(用于后续窗口计算)
) WITH (
    'connector' = 'kafka',                          -- 连接器类型:Kafka
    'topic' = 'dwd_order_topic',                    -- 关联的Kafka Topic(存储DWD层数据)
    'properties.bootstrap.servers' = 'localhost:9092',  -- Kafka集群地址
    'format' = 'json'                               -- 数据格式:JSON
);

-- 插入数据到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
-- 维表关联:用商品ID关联MySQL的product表
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,后续可以用于其他分析(如用户行为分析)。
(3)创建DWS层表:实时聚合计算

DWS层是实时数据仓库的“出口”,负责生成汇总数据,供可视化工具查询。本次实战中,我们创建两张DWS层表:

① 每小时销售额与订单量(dws_hourly_sales
-- 创建ClickHouse表(存储每小时销售额汇总数据)
CREATE TABLE 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',                     -- 连接器类型:ClickHouse
    'url' = 'clickhouse://localhost:8123',          -- ClickHouse地址
    'database-name' = 'shop',                       -- 数据库名称
    'table-name' = 'dws_hourly_sales',              -- 表名称
    'username' = 'default',                         -- ClickHouse用户名(默认:default)
    'password' = '',                                -- ClickHouse密码(默认空)
    'sink.batch-size' = '1000',                     -- 批量写入大小(1000条/批)
    'sink.flush-interval' = '1000'                  -- 刷新间隔(1秒)
);

-- 插入数据到DWS层(每小时汇总销售额与订单量)
INSERT INTO dws_hourly_sales
SELECT
    -- 将窗口开始时间格式化为“yyyy-MM-dd HH:00”(如2024-05-01 14:00)
    DATE_FORMAT(window_start, 'yyyy-MM-dd HH:00') AS hour,
    SUM(amount) AS total_sales,                     -- 总销售额(求和)
    COUNT(DISTINCT order_id) AS total_orders,       -- 总订单量(去重计数)
    window_start,                                    -- 窗口开始时间
    window_end                                       -- 窗口结束时间
FROM TABLE(
    -- 滚动窗口(Tumbling Window):每1小时一个窗口,基于事件时间
    TUMBLE(TABLE dwd_order, DESCRIPTOR(event_time), INTERVAL '1' HOUR)
)
-- 按窗口开始时间和结束时间分组(每个窗口生成一条汇总数据)
GROUP BY window_start, window_end;
② 每小时TOP 10商品(dws_hourly_top_products
-- 创建ClickHouse表(存储每小时TOP 10商品数据)
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)        -- 窗口结束时间
) WITH (
    'connector' = 'clickhouse',                     -- 连接器类型:ClickHouse
    'url' = 'clickhouse://localhost:8123',          -- ClickHouse地址
    'database-name' = 'shop',                       -- 数据库名称
    'table-name' = 'dws_hourly_top_products',       -- 表名称
    'username' = 'default',                         -- ClickHouse用户名(默认:default)
    'password' = '',                                -- ClickHouse密码(默认空)
    'sink.batch-size' = '1000',                     -- 批量写入大小(1000条/批)
    'sink.flush-interval' = '1000'                  -- 刷新间隔(1秒)
);

-- 插入数据到DWS层(每小时TOP 10商品)
INSERT INTO dws_hourly_top_products
SELECT
    DATE_FORMAT(window_start, 'yyyy-MM-dd HH:00') AS hour,  -- 小时
    product_id,                                        -- 商品ID
    product_name,                                       -- 商品名称
    sales,                                              -- 商品销售额
    rank,                                               -- 排名
    window_start,                                        -- 窗口开始时间
    window_end                                           -- 窗口结束时间
FROM (
    -- 子查询:计算每个商品在每个窗口的销售额,并排名
    SELECT
        product_id,
        product_name,
        SUM(amount) AS sales,                          -- 商品销售额(求和)
        window_start,
        window_end,
        -- 排名函数:按销售额降序排列(每个窗口内排名)
        RANK() OVER (PARTITION BY window_start ORDER BY SUM(amount) DESC) AS rank
    FROM TABLE(
        -- 滚动窗口(Tumbling Window):每1小时一个窗口,基于事件时间
        TUMBLE(TABLE dwd_order, DESCRIPTOR(event_time), INTERVAL '1' HOUR)
    )
    -- 按商品ID、商品名称、窗口开始时间、窗口结束时间分组
    GROUP BY product_id, product_name, window_start, window_end
)
-- 过滤条件:只取排名前10的商品
WHERE rank <= 10;

说明

  • 滚动窗口(Tumbling Window):每1小时一个窗口,窗口之间不重叠(如14:00-15:00、15:00-16:00),适合固定时间间隔的汇总;
  • RANK()函数:用于计算每个商品在窗口内的排名(降序),PARTITION BY window_start表示按窗口分组,ORDER BY SUM(amount) DESC表示按销售额降序排列;
  • ClickHouse存储:ClickHouse的列式存储和向量查询引擎,使得汇总数据的查询速度非常快(秒级),适合实时展示。
3. 步骤3:可视化展示——Grafana连接ClickHouse

最后,我们用Grafana连接ClickHouse,将DWS层的汇总数据展示为实时 dashboard。

(1)配置Grafana数据源
  1. 登录Grafana(默认地址:http://localhost:3000,用户名/密码:admin/admin);
  2. 点击左侧菜单栏的“Configuration”→“Data Sources”;
  3. 点击“Add data source”,选择“ClickHouse”;
  4. 配置ClickHouse连接信息:
    • URL:http://localhost:8123(ClickHouse的HTTP端口);
    • Database:shop(数据库名称);
    • User:default(默认用户名);
    • Password:空(默认密码);
  5. 点击“Save & Test”,验证连接是否成功。
(2)创建Dashboard
  1. 点击左侧菜单栏的“Create”→“Dashboard”;
  2. 点击“Add panel”,选择“Graph”(折线图,展示销售额趋势);
  3. 在“Query”标签页,选择之前配置的ClickHouse数据源,输入查询语句:
    SELECT
        hour AS "时间",
        total_sales AS "销售额"
    FROM dws_hourly_sales
    ORDER BY window_start DESC
    LIMIT 24;  -- 展示最近24小时的数据
    
  4. 调整图表设置(如标题、坐标轴标签、颜色);
  5. 重复步骤2-4,创建其他面板:
    • 柱状图:展示每小时订单量(查询dws_hourly_salestotal_orders字段);
    • 表格:展示每小时TOP 10商品(查询dws_hourly_top_productsproduct_namesalesrank字段)。
(3)最终效果

外链图片转存失败,源站可能有防盗链机制,建议将图片保存下来直接上传
(注:以上为示意图,实际效果根据数据不同而变化)

总结与扩展

1. 总结:核心步骤回顾

本次实战从需求分析到最终展示,完成了一个端到端的实时数据仓库构建,核心步骤如下:

  1. 需求分析:明确要解决的问题(实时销售额监控);
  2. 数据建模:设计ODS、DWD、DWS三层模型,明确每层的作用;
  3. 数据Pipeline实现
    • 用Python模拟订单数据,发送到Kafka的ODS层;
    • 用Flink SQL处理数据,完成DWD层清洗和DWS层聚合;
    • 将DWS层数据存储到ClickHouse,供可视化查询;
  4. 可视化展示:用Grafana连接ClickHouse,创建实时 dashboard。

2. 常见问题(FAQ)

(1)Kafka的Partition设置多少合适?

建议Partition数等于或大于Flink的并行度(如Flink的并行度为4,则Partition数设置为4或8)。这样可以提高Flink的并行处理能力,避免数据倾斜。

(2)Flink的Checkpoint怎么设置?

Checkpoint用于保证数据一致性(精确一次语义),建议设置:

  • Checkpoint间隔:1-5分钟(根据业务延迟要求调整);
  • 模式:EXACTLY_ONCE(精确一次);
  • 状态后端:RocksDB(适合大状态场景)。

在Flink SQL中,可以通过以下配置设置Checkpoint:

SET 'execution.checkpointing.interval' = '1min';
SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
SET 'state.backend' = 'rocksdb';
(3)窗口延迟怎么办?

如果数据延迟超过水印设置的时间(如5秒),可以调整水印的延迟时间:

WATERMARK FOR event_time AS event_time - INTERVAL '10' SECOND;  -- 允许10秒延迟

另外,可以用**侧输出流(Side Output)**处理迟到的数据,将迟到的数据写入单独的存储(如Hive),供后续离线分析。

3. 下一步扩展方向

(1)加入实时维度关联

比如关联用户维度表(user_dim),统计不同性别、年龄的用户销售额:

-- 创建用户维度表(MySQL)
CREATE TABLE user_dim (
    user_id INT PRIMARY KEY,
    gender STRING,
    age INT
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:mysql://localhost:3306/shop',
    'table-name' = 'user',
    'username' = 'root',
    'password' = 'root'
);

-- 关联用户维度表,统计不同性别的销售额
INSERT INTO dws_hourly_sales_by_gender
SELECT
    DATE_FORMAT(window_start, 'yyyy-MM-dd HH:00') AS hour,
    u.gender,
    SUM(amount) AS total_sales
FROM dwd_order o
JOIN user_dim FOR SYSTEM_TIME AS OF o.proc_time u 
ON o.user_id = u.user_id
GROUP BY window_start, window_end, u.gender;
(2)优化性能
  • 状态管理:用RocksDB作为状态后端,处理大状态(如长时间窗口的聚合);
  • 并行度调整:根据数据量调整Flink的并行度(如SET 'parallelism.default' = '8';);
  • 数据压缩:启用Kafka的数据压缩(如snappy),减少网络传输量。
(3)加入监控报警

用Prometheus采集Flink的Metrics(如job_manager.job.statustask_manager.num_records_in),用Alertmanager设置报警规则(如销售额突然下降50%时,发送邮件通知运维人员)。

结语

实时数据仓库是企业实现“数据驱动决策”的关键基础设施,而Kafka+Flink的组合,凭借其高吞吐量、低延迟、精确一次语义等特性,成为构建实时数据仓库的首选方案。

本文通过一个具体的电商实战案例,详细讲解了实时数据仓库的建模思路和实现步骤。希望读者能够动手实践,根据自己的业务需求调整模型,构建属于自己的实时数据仓库。

如果您有任何问题或建议,欢迎在评论区留言,我们一起讨论!

参考资料

  • Kafka官方文档:https://kafka.apache.org/documentation/
  • Flink官方文档:https://flink.apache.org/docs/stable/
  • ClickHouse官方文档:https://clickhouse.com/docs/en/
  • Grafana官方文档:https://grafana.com/docs/
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐