从零到一:基于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以上,保证基本的冗余和高可用。

核心组件清单与版本:

  1. Apache Flink 1.17.1:我们选择的计算引擎核心。1.17版本在流批一体SQL的成熟度、与数据湖集成的易用性上都有显著提升。
  2. Apache Kafka 3.4.0:作为实时数据总线。选择3.x版本是为了更好的性能和对Kraft模式(去ZooKeeper)的支持,简化部署。
  3. Apache Iceberg 1.3.0:作为统一的数据湖存储格式。它提供了ACID事务、隐藏分区、模式演化等企业级特性,与Flink集成度极高。
  4. Flink CDC 2.4+:用于捕获MySQL变更数据的神器。我们使用 flink-sql-connector-mysql-cdc,它能直接以SQL的方式定义CDC源表,极大简化开发。
  5. 实时OLAP引擎(二选一)
    • Apache Doris 1.2.4:国产优秀项目,对宽表聚合查询和标准SQL支持非常好,运维相对简单。
    • ClickHouse 22.8+:在复杂聚合查询和单表性能上表现极致,但分布式Join和更新操作是其弱项。
  6. 监控与调度: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增量读取。这保证了数据不丢失。对于大表,可以考虑使用initialschema_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)。这样做的好处是:

  1. 数据可回溯:任何时间点的聚合结果都被保存下来,可以随时查询历史任意分钟的数据。
  2. 支持批处理:如果需要重新计算历史某天的数据(比如维度表变更),可以直接以Iceberg表为源启动一个批处理Flink作业,而无需从Kafka重放。
  3. 统一服务层:下游的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分区策略不当引发的查询慢问题。我的经验是,先让管道跑通,再针对性能瓶颈逐个击破,同时建立完善的监控体系,这样才能让实时数仓真正稳定、可靠地驱动业务决策。

更多推荐