Doris 与 Flink 整合实战:构建实时计算分析平台
引言
痛点引入:实时数据处理的“最后一公里”难题
在数字化时代,实时性已经成为企业竞争的核心壁垒——电商需要实时监控订单波动,物流需要实时追踪包裹位置,游戏需要实时统计玩家行为。但很多团队在搭建实时系统时,都会遇到一个共性问题:
“计算很快,但分析很慢”
比如,你用Flink实时处理了用户点击流,计算出每5分钟的UV(独立访客),但结果存到了HBase或MySQL里。当业务人员想基于这些数据做实时Dashboard时,却发现:
- HBase不支持复杂SQL查询,得写MapReduce才能统计;
- MySQL面对千万级数据的聚合查询(比如“按省份分组查UVTop10”),延迟能高达几十秒;
- 若要做多维分析(比如“同时按时间、地域、商品分类查询”),要么得预计算大量中间表,要么得接受“查一次等半分钟”的尴尬。
这就是实时数据处理的“最后一公里”:实时计算的结果无法被高效分析。而解决这个问题的关键,在于找到一个能衔接“实时计算”与“实时分析”的桥梁。
解决方案概述:Flink + Doris,让数据“算完就能用”
Doris(Apache Doris,原百度Palo)是一款MPP架构的实时分析型数据库,主打“高吞吐写入、低延迟查询、复杂SQL支持”;而Flink是业内公认的实时计算引擎天花板,擅长处理流数据的低延迟计算。
当两者结合时,能形成一条端到端的实时数据链路:
这个方案的核心优势是:
- 低延迟:Flink的计算延迟在毫秒级,Doris的查询延迟在毫秒到秒级;
- 高吞吐:Doris支持每秒十万级别的写入,Flink支持百万级别的并发处理;
- 易扩展:两者都是分布式架构,能通过增加节点线性扩展能力;
- 简化链路:无需中间存储(比如HBase/MySQL),直接从计算到分析,减少数据移动。
最终效果展示
我们将通过实战搭建一个电商实时UV分析系统,最终实现:
- Flink实时计算每5分钟的UV(独立访客);
- 结果实时写入Doris;
- 用Doris SQL查询“近1小时各时间段的UV趋势”;
- 用Grafana展示实时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.0–fe–x86_64
container_name: doris–fe
ports:
– "8030:8030" # FE的HTTP端口(用于Flink连接)
– "9030:9030" # FE的MySQL协议端口(用于SQL客户端连接)
environment:
– FE_SERVERS=doris–fe1:10.0.2.15:9010
– FE_ID=1
volumes:
– ./doris–fe/data:/opt/apache–doris/fe/doris–meta
– ./doris–fe/log:/opt/apache–doris/fe/log
doris-be:
image: apache/doris:2.1.0–be–x86_64
container_name: doris–be
ports:
– "8040:8040" # BE的HTTP端口
depends_on:
– doris–fe
environment:
– BE_ADDR=doris–be1:10.0.2.16:9060
– FE_SERVERS=doris–fe:10.0.2.15:9030
volumes:
– ./doris–be/data:/opt/apache–doris/be/storage
– ./doris–be/log:/opt/apache–doris/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,可以看到类似以下的结果:
| 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;
若看到类似以下结果,说明写入成功:
| 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,安装步骤:
5.3.2 配置Doris数据源
在Grafana中添加数据源:
- Name:数据源名称(比如Doris);
- FENodes:Doris FE的HTTP地址(比如http://localhost:8030);
- Username:Doris用户名(root);
- Password:Doris密码(空);
5.3.3 制作实时UV Dashboard
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;
原理深入:Flink + Doris的“实时”是如何实现的?
为了让你不仅“会用”,更“懂原理”,我们深入讲解Flink Doris Connector的工作机制和Exactly-Once语义的实现。
1. Flink Doris Connector的工作原理
Flink Doris Connector基于Doris的Stream Load协议实现数据写入。Stream Load是Doris的高吞吐写入接口,支持每秒十万级别的数据写入,且延迟在毫秒级。
Connector的工作流程:
2. Exactly-Once语义的实现
Exactly-Once是实时系统的关键需求——数据不丢不重。Flink + Doris的Exactly-Once通过以下机制实现:
- 准备阶段: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整合的核心流程:
常见问题(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;
下一步:构建更完整的实时数据平台
本文的实战是一个基础的实时分析系统,你可以在此基础上扩展更复杂的场景:
结语
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/
最后,祝愿你在实时数据的世界里,越走越远!🚀




