Flink 1.13 连接 Doris 实战:从DataStream到SQL,两种方式读写数据保姆级教程
·
Flink 1.13 与 Doris 深度集成实战:双引擎数据处理的进阶指南
在实时数据处理的战场上,Apache Flink 以其卓越的流批一体能力成为事实标准,而 Apache Doris 则凭借其极速的分析性能在 OLAP 领域崭露头角。当这两个系统相遇时,会擦出怎样的火花?本文将带您深入探索 Flink 1.13 与 Doris 的深度集成方案,从底层原理到生产实践,为您呈现一份全方位的技术指南。
1. 环境准备与核心原理
在开始编码之前,我们需要先搭建好开发环境并理解两个系统协同工作的基本原理。Flink-Doris Connector 的实现基于 Doris 的 Stream Load 协议,这是一种高效的批量数据导入机制。
1.1 依赖配置与版本适配
在项目的 pom.xml 中添加以下依赖(以 Maven 为例):
<dependency>
<groupId>org.apache.doris</groupId>
<artifactId>flink-doris-connector-1.13_2.12</artifactId>
<version>1.0.3</version>
</dependency>
版本匹配至关重要,以下是常见的兼容性对照表:
| Flink 版本 | Scala 版本 | Connector 版本 |
|---|---|---|
| 1.13.x | 2.12 | 1.0.x |
| 1.14.x | 2.12 | 1.1.x |
| 1.15.x | 2.12 | 1.2.x |
提示:生产环境中建议使用相同 minor 版本的最新 patch 版本,以获得最稳定的表现。
1.2 Doris 表设计要点
在 Doris 中创建测试表时,需要特别注意与 Flink 的兼容性:
CREATE TABLE flink_demo (
user_id BIGINT,
event_time DATETIME,
device_id VARCHAR(128),
metric_value DECIMAL(18,4)
)
UNIQUE KEY(user_id, event_time)
DISTRIBUTED BY HASH(user_id) BUCKETS 16
PROPERTIES (
"replication_num" = "3",
"storage_medium" = "SSD"
);
关键设计考量:
- Unique Key 模型 :确保数据可更新
- 合理分桶数 :建议是节点数的 3-5 倍
- 数据类型映射 :特别注意 DECIMAL 精度和 VARCHAR 长度
2. DataStream API 深度集成
对于习惯使用 Flink 底层 API 的开发者,DataStream 方式提供了最灵活的控制能力。
2.1 高效读取 Doris 数据
以下是一个完整的 Doris 数据源读取示例:
Properties readProps = new Properties();
readProps.setProperty("fenodes", "doris-fe:8030");
readProps.setProperty("username", "flink_user");
readProps.setProperty("password", "password123");
readProps.setProperty("table.identifier", "db1.user_behavior");
readProps.setProperty("read.fields", "user_id,event_time,action_type");
DorisSourceFunction<List<?>> source = DorisSourceFunction.<List<?>>builder()
.setDorisOptions(new DorisStreamOptions(readProps))
.setDeserializer(new SimpleListDeserializationSchema())
.build();
DataStream<List<?>> dorisStream = env.addSource(source)
.setParallelism(2)
.name("DorisSource");
性能调优参数:
| 参数名 | 默认值 | 建议值 | 说明 |
|---|---|---|---|
| doris.request.query.timeout | 3600 | 1800 | 查询超时时间(秒) |
| doris.request.connect.timeout | 30000 | 10000 | 连接超时时间(毫秒) |
| doris.request.read.timeout | 30000 | 30000 | 读取超时时间(毫秒) |
| doris.request.retry.times | 3 | 5 | 查询失败重试次数 |
| doris.batch.size | 1024 | 4096 | 单次读取批大小 |
2.2 数据写入的两种模式
JSON 格式写入
DorisSink<String> sink = DorisSink.<String>builder()
.setDorisOptions(DorisOptions.builder()
.setFenodes("doris-fe:8030")
.setUsername("flink_user")
.setPassword("password123")
.setTableIdentifier("db1.user_actions")
.build())
.setDorisExecutionOptions(DorisExecutionOptions.builder()
.setBatchIntervalMs(5000L)
.setBatchSize(1024 * 1024)
.setMaxRetries(3)
.setStreamLoadProp(ImmutableMap.of(
"format", "json",
"strip_outer_array", "true"
))
.build())
.setSerializer((value, out) -> {
out.write(value.getBytes(StandardCharsets.UTF_8));
})
.build();
stream.addSink(sink).name("DorisSink");
RowData 格式写入(更高性能)
LogicalType[] fieldTypes = {
new BigIntType(),
new TimestampType(),
new VarCharType(128),
new DecimalType(18, 4)
};
String[] fieldNames = {"user_id", "event_time", "device_id", "metric_value"};
DorisSink<RowData> sink = DorisSink.<RowData>builder()
.setDorisOptions(DorisOptions.builder()
.setFenodes("doris-fe:8030")
.setUsername("flink_user")
.setPassword("password123")
.setTableIdentifier("db1.metrics")
.build())
.setDorisExecutionOptions(DorisExecutionOptions.builder()
.setBatchIntervalMs(2000L)
.setEnableDelete(false)
.setMaxRetries(3)
.build())
.setRowDataSerializer(fieldNames, fieldTypes)
.build();
3. SQL API 优雅集成
对于偏好声明式编程的团队,Flink SQL 提供了更简洁的集成方式。
3.1 创建 Doris 表映射
CREATE TABLE doris_user_actions (
`user_id` BIGINT,
`action_time` TIMESTAMP(3),
`page_id` VARCHAR(100),
`action_type` VARCHAR(20)
) WITH (
'connector' = 'doris',
'fenodes' = 'doris-fe:8030',
'table.identifier' = 'db1.user_actions',
'username' = 'flink_user',
'password' = 'password123',
'sink.batch.size' = '1000',
'sink.batch.interval' = '1s'
);
3.2 复杂数据处理示例
-- 实时维表关联
INSERT INTO doris_output
SELECT
e.user_id,
e.event_time,
u.user_name,
u.user_level,
COUNT(*) AS event_count
FROM kafka_events AS e
JOIN doris_users FOR SYSTEM_TIME AS OF e.proc_time AS u
ON e.user_id = u.user_id
GROUP BY
e.user_id, e.event_time, u.user_name, u.user_level;
-- 时间窗口聚合
INSERT INTO doris_metrics
SELECT
window_start,
window_end,
device_type,
COUNT(DISTINCT user_id) AS uv,
SUM(click_count) AS total_clicks
FROM TABLE(
TUMBLE(TABLE user_clicks, DESCRIPTOR(event_time), INTERVAL '5' MINUTES)
)
GROUP BY window_start, window_end, device_type;
4. 生产环境最佳实践
4.1 性能调优指南
写入优化:
- 调整
sink.batch.size(默认1000) 和sink.batch.interval(默认1s) - 对于高吞吐场景,建议值:
- batch.size: 5000-10000
- batch.interval: 2-5s
读取优化:
-- 启用分区裁剪和谓词下推
SELECT * FROM doris_table
WHERE dt='2023-01-01' AND user_id > 1000
WITH ('scan.partition.column'='dt', 'scan.partition.values'='2023-01-01');
4.2 容错与监控
关键监控指标:
doris.sink.batch.count:每批次写入记录数doris.sink.duration:写入耗时doris.sink.success:成功批次计数doris.sink.failed:失败批次计数
异常处理策略:
DorisExecutionOptions.builder()
.setMaxRetries(5)
.setRetryBackoffMultiplierMs(1000)
.setEnable2PC(true) // 启用两阶段提交
.setIgnoreUpdateBefore(true) // 忽略更新前状态
.build()
4.3 数据类型映射陷阱
实际工作中需特别注意的类型转换:
| Doris 类型 | Flink SQL 类型 | 注意事项 |
|---|---|---|
| LARGEINT | STRING | 可能丢失精度 |
| DECIMALV2 | DECIMAL(p,s) | 需确保精度一致 |
| DATETIME | TIMESTAMP(3) | 时区处理需一致 |
| HLL | 不支持 | 需在Doris侧预处理 |
| ARRAY/JSON | STRING | 需要额外序列化处理 |
5. 典型应用场景解析
5.1 实时数仓构建
-- 分层建模示例
CREATE TABLE ods_events (
-- 字段定义
) WITH ('connector'='kafka', ...);
CREATE TABLE dwd_actions (
-- 字段定义
) WITH ('connector'='doris', ...);
-- ODS到DWD处理
INSERT INTO dwd_actions
SELECT
user_id,
event_time,
-- 维度解析
JSON_VALUE(properties, '$.page_id') AS page_id,
-- 指标计算
COUNT(*) AS pv
FROM ods_events
GROUP BY user_id, event_time, ...;
5.2 维表实时关联
// 使用Async I/O实现高效维表关联
AsyncDataStream.unorderedWait(
eventStream,
new DorisAsyncLookupFunction(
"jdbc:doris://doris-fe:9030/db1",
"users",
"user_id,user_name,gender,age",
"user_id"),
5000, // 超时时间
TimeUnit.MILLISECONDS,
100 // 最大并发请求数
);
5.3 数据湖分析加速
-- 将Hive数据加速导入Doris
INSERT INTO doris_analysis_table
SELECT * FROM hive_catalog.db1.source_table
WHERE dt='2023-01-01';
-- 定时增量同步
INSERT INTO doris_analysis_table
SELECT * FROM hive_catalog.db1.source_table
WHERE dt='2023-01-02' AND modified_time > '2023-01-01 00:00:00';
通过以上五个维度的深入探讨,我们不仅掌握了 Flink 与 Doris 集成的技术细节,更了解了如何在实际业务中发挥两者的组合优势。无论是简单的数据同步还是复杂的实时分析场景,这套技术组合都能提供令人满意的解决方案。
更多推荐



所有评论(0)