欢迎光临
我们一直在努力

Doris 与 Flink 整合实战:构建实时计算分析平台

Doris 与 Flink 整合实战:构建实时计算分析平台

引言

痛点引入:实时数据处理的“最后一公里”难题

在数字化时代,实时性已经成为企业竞争的核心壁垒——电商需要实时监控订单波动,物流需要实时追踪包裹位置,游戏需要实时统计玩家行为。但很多团队在搭建实时系统时,都会遇到一个共性问题:

“计算很快,但分析很慢”

比如,你用Flink实时处理了用户点击流,计算出每5分钟的UV(独立访客),但结果存到了HBase或MySQL里。当业务人员想基于这些数据做实时Dashboard时,却发现:

  • HBase不支持复杂SQL查询,得写MapReduce才能统计;
  • MySQL面对千万级数据的聚合查询(比如“按省份分组查UVTop10”),延迟能高达几十秒;
  • 若要做多维分析(比如“同时按时间、地域、商品分类查询”),要么得预计算大量中间表,要么得接受“查一次等半分钟”的尴尬。

这就是实时数据处理的“最后一公里”:实时计算的结果无法被高效分析。而解决这个问题的关键,在于找到一个能衔接“实时计算”与“实时分析”的桥梁。

解决方案概述:Flink + Doris,让数据“算完就能用”

Doris(Apache Doris,原百度Palo)是一款MPP架构的实时分析型数据库,主打“高吞吐写入、低延迟查询、复杂SQL支持”;而Flink是业内公认的实时计算引擎天花板,擅长处理流数据的低延迟计算。

当两者结合时,能形成一条端到端的实时数据链路:

  • 数据接入:从Kafka、Binlog、IoT设备等数据源获取实时流;
  • 实时计算:用Flink做清洗、转换、聚合(比如计算UV、订单金额、用户留存);
  • 实时写入:通过Flink Doris Connector将计算结果实时写入Doris;
  • 实时分析:用Doris的SQL做多维分析、即席查询,或对接Grafana/Tableau生成实时Dashboard。
  • 这个方案的核心优势是:

    • 低延迟:Flink的计算延迟在毫秒级,Doris的查询延迟在毫秒到秒级;
    • 高吞吐:Doris支持每秒十万级别的写入,Flink支持百万级别的并发处理;
    • 易扩展:两者都是分布式架构,能通过增加节点线性扩展能力;
    • 简化链路:无需中间存储(比如HBase/MySQL),直接从计算到分析,减少数据移动。

    最终效果展示

    我们将通过实战搭建一个电商实时UV分析系统,最终实现:

    • Flink实时计算每5分钟的UV(独立访客);
    • 结果实时写入Doris;
    • 用Doris SQL查询“近1小时各时间段的UV趋势”;
    • 用Grafana展示实时Dashboard(如下图)。

    实时UV Dashboard示例(注:此为模拟图,实际效果需自行搭建)

    准备工作

    在开始实战前,需要完成环境搭建和基础知识储备。

    1. 环境与工具清单

    工具/组件版本要求说明
    Apache Flink 1.17+ 实时计算引擎,推荐使用1.17及以上版本(支持更完善的Connector)
    Apache Doris 2.1+ 实时分析数据库,2.0+版本支持更优的Stream Load性能
    Java 1.8/11 Flink和Doris均依赖Java
    Maven/Gradle 3.6+ 构建Flink作业的依赖管理工具
    Kafka(可选) 2.8+ 模拟实际场景的数据源(若没有Kafka,可用Flink的DataGenSource替代)
    Grafana(可选) 9.0+ 可视化Dashboard工具

    2. 依赖配置(Flink作业)

    要让Flink能写入Doris,需在pom.xml(Maven)或build.gradle(Gradle)中添加Flink Doris Connector依赖:

    Maven配置

    <dependencies>
    <!– Flink Table API 基础依赖 –>
    <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-api-java-bridge</artifactId>
    <version>${flink.version}</version>
    <scope>provided</scope>
    </dependency>
    <!– Flink Doris Connector –>
    <dependency>
    <groupId>org.apache.doris</groupId>
    <artifactId>flink-doris-connector-1.17</artifactId>
    <version>1.4.0</version>
    </dependency>
    <!– 可选:Kafka Connector(若用Kafka做数据源) –>
    <dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-kafka</artifactId>
    <version>${flink.version}</version>
    </dependency>
    </dependencies>

    Gradle配置

    dependencies {
    implementation 'org.apache.flink:flink-table-api-java-bridge:${flink.version}'
    implementation 'org.apache.doris:flink-doris-connector-1.17:1.4.0'
    implementation 'org.apache.flink:flink-connector-kafka:${flink.version}'
    }

    3. 基础知识储备

    为了顺利完成实战,你需要了解以下概念:

    • Flink基础:DataStream API/Table API、Watermark(水位线)、Checkpoint( checkpoint)、窗口函数(Tumble Window/滑动窗口);
    • Doris基础:表结构(Duplicate Key/Unique Key/Aggregate Key)、分区(Partition)、分桶(Bucket)、Stream Load(流式写入协议);
    • 实时计算概念:流处理vs批处理、Exactly-Once(精确一次)语义。

    核心步骤:从0到1搭建实时分析平台

    我们将按**“环境部署→数据准备→Flink计算→写入Doris→实时分析”**的流程逐步实现。

    步骤1:快速部署Doris测试集群

    Doris支持多种部署方式(二进制、Docker、K8s),这里推荐Docker Compose快速搭建测试集群(生产环境建议用二进制部署)。

    1.1 安装Docker Compose

    确保你的机器已安装Docker和Docker Compose(若未安装,参考Docker官方文档)。

    1.2 编写Docker Compose文件

    创建docker-compose.yml,内容如下(基于Doris 2.1版本):

    version: '3'
    services:
    doris-fe:
    image: apache/doris:2.1.0fex86_64
    container_name: dorisfe
    ports:
    "8030:8030" # FE的HTTP端口(用于Flink连接)
    "9030:9030" # FE的MySQL协议端口(用于SQL客户端连接)
    environment:
    FE_SERVERS=dorisfe1:10.0.2.15:9010
    FE_ID=1
    volumes:
    ./dorisfe/data:/opt/apachedoris/fe/dorismeta
    ./dorisfe/log:/opt/apachedoris/fe/log

    doris-be:
    image: apache/doris:2.1.0bex86_64
    container_name: dorisbe
    ports:
    "8040:8040" # BE的HTTP端口
    depends_on:
    dorisfe
    environment:
    BE_ADDR=dorisbe1:10.0.2.16:9060
    FE_SERVERS=dorisfe:10.0.2.15:9030
    volumes:
    ./dorisbe/data:/opt/apachedoris/be/storage
    ./dorisbe/log:/opt/apachedoris/be/log

    1.3 启动Doris集群

    在docker-compose.yml所在目录执行:

    docker-compose up -d

    启动成功后,可通过以下方式验证:

    • 访问FE的Web UI:http://localhost:8030(默认用户名root,密码为空);
    • 用MySQL客户端连接FE:mysql -h localhost -P 9030 -u root -p(密码为空)。

    步骤2:准备实时数据源

    为了模拟真实场景,我们用Flink的DataGenSource生成模拟的电商用户行为数据(若你有Kafka集群,也可以用Kafka作为数据源)。

    2.1 定义用户行为数据结构

    用户行为数据包含以下字段:

    字段名类型说明
    user_id BIGINT 用户ID
    item_id BIGINT 商品ID
    category_id BIGINT 商品分类ID
    behavior STRING 用户行为(click/order/pay)
    ts TIMESTAMP(3) 行为发生时间
    2.2 用Flink Table API创建数据源表

    在Flink SQL客户端或代码中,创建user_behavior表(用DataGenSource生成模拟数据):

    CREATE TABLE user_behavior (
    user_id BIGINT,
    item_id BIGINT,
    category_id BIGINT,
    behavior STRING,
    ts TIMESTAMP(3),
    — 定义Watermark:处理乱序数据,延迟5秒
    WATERMARK FOR ts AS ts INTERVAL '5' SECOND
    ) WITH (
    'connector' = 'datagen', — 用DataGenSource生成数据
    'rows-per-second' = '100', — 每秒生成100条数据
    'fields.user_id.min' = '1', — user_id的最小值
    'fields.user_id.max' = '10000', — user_id的最大值
    'fields.behavior.length' = '1' — behavior字段的长度(这里用枚举值,后续过滤)
    );

    步骤3:Flink实时计算——计算每5分钟UV

    UV(独立访客)是电商最核心的实时指标之一,我们需要计算每5分钟内的独立用户数。

    3.1 窗口函数选择

    UV计算需要按时间窗口聚合,这里选择滚动窗口(Tumble Window)——每5分钟一个窗口,窗口不重叠。

    3.2 编写Flink计算SQL

    用Flink Table API编写计算逻辑:

    — 计算每5分钟的UV(仅统计click行为)
    CREATE VIEW real_time_uv_view AS
    SELECT
    TUMBLE_START(ts, INTERVAL '5' MINUTE) AS window_start, — 窗口开始时间
    TUMBLE_END(ts, INTERVAL '5' MINUTE) AS window_end, — 窗口结束时间
    COUNT(DISTINCT user_id) AS uv — 独立用户数
    FROM user_behavior
    WHERE behavior = 'click' — 仅统计点击行为
    GROUP BY TUMBLE(ts, INTERVAL '5' MINUTE); — 按时间滚动窗口分组

    3.3 验证计算逻辑

    在Flink SQL客户端执行SELECT * FROM real_time_uv_view,可以看到类似以下的结果:

    window_startwindow_enduv
    2024-05-20 10:00:00 2024-05-20 10:05:00 123
    2024-05-20 10:05:00 2024-05-20 10:10:00 456

    步骤4:整合Flink与Doris——实时写入计算结果

    接下来,我们需要将Flink计算出的UV结果实时写入Doris,这是整个链路的核心环节。

    4.1 在Doris中创建目标表

    首先,在Doris中创建用于存储UV结果的表real_time_uv。注意:

    • 选择Duplicate Key表类型(适合存储原始聚合结果,支持更新);
    • 按window_start(窗口开始时间)分区(方便按时间范围查询);
    • 按window_start分桶(均匀分布数据)。

    执行以下SQL(通过MySQL客户端连接Doris):

    CREATE DATABASE IF NOT EXISTS test; — 创建测试数据库

    USE test;

    CREATE TABLE real_time_uv (
    window_start DATETIME, — 窗口开始时间
    window_end DATETIME, — 窗口结束时间
    uv BIGINT — 独立访客数
    )
    DUPLICATE KEY (window_start, window_end) — 主键(用于去重)
    PARTITION BY RANGE(window_start) ( — 按时间分区
    START ('2024-01-01') END ('2024-12-31') EVERY (INTERVAL 1 DAY)
    )
    DISTRIBUTED BY HASH(window_start) BUCKETS 10 — 分桶(10个桶)
    PROPERTIES (
    "replication_num" = "1", — 副本数(测试环境设为1)
    "in_memory" = "true" — 内存中存储(加速查询)
    );

    4.2 在Flink中配置Doris Sink

    在Flink中创建doris_sink表,用于将计算结果写入Doris。关键配置参数说明:

    参数说明
    connector 固定为doris(指定使用Flink Doris Connector)
    fenodes Doris FE的HTTP地址(格式:FE_IP:FE_HTTP_PORT,比如localhost:8030)
    username Doris的用户名(默认root)
    password Doris的密码(默认空)
    table.identifier Doris的表标识符(格式:数据库名.表名,比如test.real_time_uv)
    sink.properties.format 写入数据的格式(支持json或csv,这里用json)
    sink.properties.strip_outer_array 去除JSON数组的外层(Flink输出的是数组,Doris需要单条JSON)
    sink.enable-delete 是否支持删除(可选,默认false)

    Flink SQL配置:

    CREATE TABLE doris_sink (
    window_start TIMESTAMP(3),
    window_end TIMESTAMP(3),
    uv BIGINT,
    PRIMARY KEY (window_start, window_end) NOT ENFORCED — 主键(与Doris表对齐)
    ) WITH (
    'connector' = 'doris',
    'fenodes' = 'localhost:8030',
    'username' = 'root',
    'password' = '',
    'table.identifier' = 'test.real_time_uv',
    'sink.properties.format' = 'json',
    'sink.properties.strip_outer_array' = 'true',
    'sink.parallelism' = '2' — 写入Doris的并行度(根据集群规模调整)
    );

    4.3 启动Flink写入任务

    执行INSERT INTO语句,将real_time_uv_view的结果写入doris_sink:

    INSERT INTO doris_sink
    SELECT * FROM real_time_uv_view;

    4.4 验证数据写入

    在Doris中执行查询,确认数据已写入:

    SELECT window_start, window_end, uv FROM test.real_time_uv ORDER BY window_start DESC LIMIT 10;

    若看到类似以下结果,说明写入成功:

    window_startwindow_enduv
    2024-05-20 10:30:00 2024-05-20 10:35:00 789
    2024-05-20 10:25:00 2024-05-20 10:30:00 654
    2024-05-20 10:20:00 2024-05-20 10:25:00 321

    步骤5:Doris实时分析——从数据到价值

    Doris的核心优势是复杂SQL的低延迟查询。我们可以基于real_time_uv表做多种实时分析。

    5.1 实时查询:近1小时的UV趋势

    查询最近1小时内,每5分钟的UV变化:

    SELECT
    window_start,
    uv
    FROM test.real_time_uv
    WHERE window_start >= NOW() INTERVAL '1' HOUR — 过滤最近1小时的数据
    ORDER BY window_start ASC;

    5.2 多维分析:按天统计UV总和

    若要统计每天的总UV,只需在window_start上做日期截断:

    SELECT
    DATE(window_start) AS dt, — 截断到日期
    SUM(uv) AS total_uv — 求和
    FROM test.real_time_uv
    GROUP BY dt
    ORDER BY dt DESC;

    5.3 实时Dashboard:用Grafana可视化

    为了让业务人员更直观地看到数据,我们用Grafana对接Doris,生成实时Dashboard。

    5.3.1 安装Grafana Doris插件

    Grafana社区提供了Doris的数据源插件grafana-doris-datasource,安装步骤:

  • 下载插件:git clone https://github.com/apache/doris/tree/master/grafana-datasource;
  • 将插件复制到Grafana的插件目录(比如/var/lib/grafana/plugins);
  • 重启Grafana:sudo systemctl restart grafana-server。
  • 5.3.2 配置Doris数据源

    在Grafana中添加数据源:

  • 点击左侧菜单栏的“Configuration”→“Data Sources”;
  • 点击“Add data source”,搜索“Doris”;
  • 配置参数:
    • Name:数据源名称(比如Doris);
    • FENodes:Doris FE的HTTP地址(比如http://localhost:8030);
    • Username:Doris用户名(root);
    • Password:Doris密码(空);
  • 点击“Save & Test”,验证连接成功。
  • 5.3.3 制作实时UV Dashboard
  • 点击左侧菜单栏的“Create”→“Dashboard”;
  • 点击“Add panel”,选择“Table”或“Time series”图表;
  • 编写SQL查询(比如“近1小时的UV趋势”):SELECT
    window_start AS "time", — Grafana需要`time`字段作为时间轴
    uv AS "UV"
    FROM test.real_time_uv
    WHERE window_start >= $__timeFrom() AND window_start <= $__timeTo() — 时间范围变量
    ORDER BY window_start ASC;
  • 调整图表样式(比如选择“Line”类型,设置刷新间隔为“5s”);
  • 保存Dashboard,完成!
  • 原理深入:Flink + Doris的“实时”是如何实现的?

    为了让你不仅“会用”,更“懂原理”,我们深入讲解Flink Doris Connector的工作机制和Exactly-Once语义的实现。

    1. Flink Doris Connector的工作原理

    Flink Doris Connector基于Doris的Stream Load协议实现数据写入。Stream Load是Doris的高吞吐写入接口,支持每秒十万级别的数据写入,且延迟在毫秒级。

    Connector的工作流程:

  • 数据分片:Flink将数据分成多个分片(由sink.parallelism控制并行度);
  • 数据缓存:每个分片将数据缓存到内存中(默认缓存大小为1MB,可通过sink.batch.size调整);
  • 批量写入:当缓存达到阈值或超时(sink.batch.interval)时,将数据打包成JSON数组,通过Stream Load协议发送到Doris;
  • 结果确认:Doris返回写入结果,Flink根据结果确认是否提交Checkpoint(保证Exactly-Once)。
  • 2. Exactly-Once语义的实现

    Exactly-Once是实时系统的关键需求——数据不丢不重。Flink + Doris的Exactly-Once通过以下机制实现:

  • Flink Checkpoint:Flink定期生成Checkpoint,记录当前的处理位置;
  • Doris事务:Doris的Stream Load支持事务(通过label参数唯一标识一个事务);
  • 两阶段提交(2PC):
    • 准备阶段:Flink将数据写入Doris的临时目录,不提交事务;
    • 提交阶段:Flink完成Checkpoint后,通知Doris提交事务;
    • 回滚阶段:若Flink任务失败,从最近的Checkpoint恢复,重新写入数据,Doris会自动忽略重复的label。
  • 要开启Exactly-Once,需要在Flink中配置Checkpoint:

    StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
    env.enableCheckpointing(60000); // 每60秒做一次Checkpoint
    env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 开启Exactly-Once

    3. 性能优化技巧

    为了让系统达到最佳性能,你可以调整以下参数:

    参数优化方向
    sink.batch.size 增大缓存大小(比如从1MB调整到10MB),减少Stream Load的次数,降低延迟
    sink.batch.interval 增大超时时间(比如从1s调整到5s),合并更多数据写入
    sink.parallelism 增加并行度(比如从2调整到10),提高写入吞吐
    doris.be.number 增加Doris BE的数量(生产环境建议至少3个),提高查询性能
    doris.table.buckets 调整分桶数(比如从10调整到20),均匀分布数据

    总结与扩展

    回顾要点

    通过本文的实战,你已经掌握了Flink + Doris整合的核心流程:

  • 部署环境:用Docker快速启动Doris集群,配置Flink依赖;
  • 准备数据:用DataGenSource生成模拟用户行为数据;
  • 实时计算:用Flink的窗口函数计算每5分钟的UV;
  • 写入Doris:配置Flink Doris Connector,将结果实时写入Doris;
  • 实时分析:用Doris SQL做多维查询,用Grafana可视化。
  • 常见问题(FAQ)

    Q1:Flink任务启动失败,提示“无法连接到Doris FE”?
    • 检查fenodes参数是否正确(比如localhost:8030,注意不是MySQL端口9030);
    • 检查Doris FE是否在运行(docker ps看doris-fe容器状态);
    • 检查防火墙是否开放了8030端口(sudo ufw allow 8030)。
    Q2:写入Doris的数据查询不到?
    • 检查Flink任务是否在运行(flink list看任务状态);
    • 检查Doris的分区是否正确(比如window_start是否在分区范围内);
    • 检查Flink的watermark设置是否合理(比如延迟5秒,是否数据还没到窗口结束时间)。
    Q3:查询延迟高怎么办?
    • 调整Doris表的in_memory参数为true(内存中存储);
    • 增加Doris BE的数量(提高查询并行度);
    • 为高频查询创建Rollup表(预聚合,加速查询):CREATE ROLLUP real_time_uv_rollup (
      window_start,
      uv
      )
      FROM real_time_uv
      KEY (window_start)
      DISTRIBUTED BY HASH(window_start) BUCKETS 10;

    下一步:构建更完整的实时数据平台

    本文的实战是一个基础的实时分析系统,你可以在此基础上扩展更复杂的场景:

  • 对接Kafka:用Kafka作为数据源(更贴近实际生产);
  • 多流Join:用Flink做双流Join(比如用户行为流与商品信息流Join);
  • 实时数仓:结合Hudi/Deltalake构建实时数仓,Doris作为数仓的分析层;
  • 告警系统:用Flink的CEP(复杂事件处理)检测异常数据,触发告警(比如UV骤降)。
  • 结语

    Flink + Doris的整合,解决了实时数据处理的“最后一公里”问题——让计算结果能被快速分析。无论是电商的实时监控、物流的实时追踪,还是游戏的实时统计,这个组合都能胜任。

    通过本文的实战,你已经掌握了搭建实时分析平台的核心能力。接下来,不妨尝试将其应用到你的实际项目中,让数据真正“活”起来!

    如果在实践中遇到问题,欢迎在评论区交流,或参考以下资源:

    • Doris官方文档:https://doris.apache.org/
    • Flink官方文档:https://flink.apache.org/
    • Flink Doris Connector文档:https://doris.apache.org/docs/dev/ecosystem/flink-doris-connector/

    最后,祝愿你在实时数据的世界里,越走越远!🚀

    赞(0)
    未经允许不得转载:171主机测评 » Doris 与 Flink 整合实战:构建实时计算分析平台
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址