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 集成的技术细节,更了解了如何在实际业务中发挥两者的组合优势。无论是简单的数据同步还是复杂的实时分析场景,这套技术组合都能提供令人满意的解决方案。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐