Flink + Doris 黄金搭档:手把手搭建毫秒级延迟的实时数仓

最近两年,我身边不少做数据平台的朋友都在经历一场“架构瘦身”运动。过去那种动辄十几个组件堆叠起来的大数据平台,维护起来像在走钢丝,一个组件升级就可能引发一连串的兼容性问题。更头疼的是,业务方对数据时效性的要求越来越高,从T+1到小时级,再到分钟级,现在甚至要求秒级可见。传统的Lambda架构,同一份业务逻辑要写两套代码(实时和离线),开发效率低不说,数据一致性还经常出问题。

正是在这种背景下,FlinkApache Doris 的组合开始频繁出现在技术选型的讨论桌上。Flink作为流计算的标杆,其CDC(变更数据捕获)能力已经非常成熟,能精准捕获数据库的每一行变化。而Doris,这个近几年在开源社区异常活跃的MPP分析型数据库,以其极致的查询性能、MySQL协议兼容性和运维简单著称。把它们俩放在一起,目标很明确:用最简洁的架构,构建一个从数据产生到可分析仅需秒级、查询响应在毫秒级的实时数据仓库。

这篇文章,我就结合自己最近在几个项目中落地的经验,抛开那些宽泛的概念,直接切入实战。我会带你一步步搭建一个从MySQL业务库到Doris的完整实时同步管道,处理生产环境中常见的分库分表合并、Schema变更等棘手问题,并提供一个开箱即用的docker-compose实验环境。无论你是想快速验证方案,还是为团队寻找一个可落地的生产级解决方案,希望这些内容都能给你带来直接的帮助。

1. 环境准备与核心组件解析

在动手搭建之前,我们得先搞清楚手里的“工具”到底能干什么。很多人对Flink的印象还停留在流计算引擎,对Doris的印象可能只是个“查得快的数据库”。其实,它们组合起来的能力远超想象。

Flink CDC 在这里扮演的是“数据搬运工+初级加工者”的角色。它不仅仅是简单地读取数据库的binlog。以Flink CDC for MySQL为例,它提供了几种核心模式:

  • 初始快照读取:全量拉取历史数据,这是启动同步的基础。
  • 增量流读取:持续监听binlog,捕获增、删、改操作。
  • Exactly-Once语义:确保数据在发生故障恢复后不丢不重,这对财务、订单等核心数据至关重要。

Apache Doris 则是一个“全能型仓库管理员”。它的核心优势在于,一份存储同时支撑多种负载。这意味着,你同步进来的数据,既可以用于高并发的实时仪表盘查询(点查询),也可以用于复杂的即席分析(Ad-Hoc),甚至可以直接在库内进行轻量级的ETL加工,无需再把数据导出到Spark或Hive。这种“离在线一体化”的特性,是简化架构的关键。

为了快速开始,我准备了一个docker-compose文件,它包含了MySQL(模拟业务库)、Flink(包含CDC连接器)、Doris(FE和BE)以及一个简易的数据生成器。你只需要有Docker环境,就能在本地完整复现整个流程。

# 1. 克隆实验环境仓库(此处为示意,实际需提供有效仓库)
# git clone https://github.com/your-repo/flink-doris-cdc-demo.git
# cd flink-doris-cdc-demo

# 2. 启动所有服务
docker-compose up -d

# 3. 查看服务状态
docker-compose ps

这个环境启动后,你会拥有:

  • MySQL: 3306端口,内置了一个简单的 user_db 库和 user_* 表(模拟分表)。
  • Flink JobManager & TaskManager: 8081端口(Web UI),已预置了Flink CDC和Doris Connector的JAR包。
  • Doris FE: 8030 (HTTP), 9030 (MySQL协议)
  • Doris BE: 8040 (HTTP)
  • Data Generator: 一个不断向MySQL插入、更新数据的小程序,用于模拟线上流量。

注意:生产环境的部署需要考虑高可用。Flink on YARN/K8s,Doris多FE/BE节点部署是标准做法。这个实验环境旨在让你无负担地快速理解流程。

2. 从MySQL到Doris:CDC实时同步实战

环境就绪,现在我们开始构建第一条实时管道。假设我们的业务库 user_db 中有一张核心表 user_info,我们需要将它实时同步到Doris的 doris_db 库中。

首先,在Doris中创建目标库和表。Doris建表时,选择正确的表模型对后续性能影响巨大。对于CDC同步这种需要处理更新的场景,我们通常使用 Unique Key 模型(Merge-on-Write)

-- 在Doris FE(通过MySQL客户端连接9030端口)执行
CREATE DATABASE IF NOT EXISTS doris_db;

USE doris_db;

CREATE TABLE IF NOT EXISTS user_info_sink (
    `user_id` BIGINT NOT NULL,
    `username` VARCHAR(50),
    `email` VARCHAR(100),
    `created_at` DATETIME,
    `updated_at` DATETIME,
    `is_deleted` TINYINT DEFAULT '0'
)
UNIQUE KEY(`user_id`) -- 指定唯一键,用于数据更新和删除
DISTRIBUTED BY HASH(`user_id`) BUCKETS 10 -- 分桶,影响数据分布和并行度
PROPERTIES (
    "replication_num" = "1", -- 实验环境副本数为1,生产环境通常为3
    "enable_unique_key_merge_on_write" = "true" -- 启用Merge-on-Write,点查性能更优
);

接下来是重头戏:编写Flink SQL作业。Flink CDC Connector 2.3+版本之后,使用起来非常直观。以下是一个完整的作业示例,它完成了数据读取、转换和写入。

-- 在Flink SQL Client中执行,或保存为SQL文件通过REST API提交
-- 1. 创建MySQL CDC源表
CREATE TABLE mysql_user_source (
    `id` BIGINT,
    `name` VARCHAR(255),
    `email` VARCHAR(255),
    `create_time` TIMESTAMP(3),
    `update_time` TIMESTAMP(3),
    `deleted` BOOLEAN,
    PRIMARY KEY (`id`) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'mysql',
    'port' = '3306',
    'username' = 'root',
    'password' = '123456',
    'database-name' = 'user_db',
    'table-name' = 'user_info',
    'server-time-zone' = 'Asia/Shanghai',
    'scan.startup.mode' = 'initial' -- 首次启动时先做全量快照,再接增量
);

-- 2. 创建Doris结果表
CREATE TABLE doris_user_sink (
    `user_id` BIGINT,
    `username` VARCHAR(50),
    `email` VARCHAR(100),
    `created_at` DATETIME,
    `updated_at` DATETIME,
    `is_deleted` TINYINT
) WITH (
    'connector' = 'doris',
    'fenodes' = 'doris-fe:8030', -- Doris FE的HTTP地址
    'table.identifier' = 'doris_db.user_info_sink',
    'username' = 'root',
    'password' = '',
    'sink.properties.format' = 'json',
    'sink.properties.read_json_by_line' = 'true', -- 每行一个JSON,是推荐格式
    'sink.buffer-flush.max-rows' = '5000', -- 批量写入参数,平衡延迟和吞吐
    'sink.buffer-flush.interval' = '2s'
);

-- 3. 将数据从MySQL写入Doris
INSERT INTO doris_user_sink
SELECT
    id AS user_id,
    name AS username,
    email,
    CAST(create_time AS TIMESTAMP) AS created_at, -- 类型转换示例
    CAST(update_time AS TIMESTAMP) AS updated_at,
    CAST(deleted AS TINYINT) AS is_deleted
FROM mysql_user_source;

提交这个作业后,Flink会先全量拉取 user_info 表的历史数据写入Doris,然后持续监听binlog,将增量变更(INSERT/UPDATE/DELETE)实时地应用到Doris表中。你可以在Flink Web UI的 Running Jobs 中看到这个作业,监控其吞吐量和延迟。

关键参数调优

  • 'scan.startup.mode':除了 initial,还有 latest-offset(仅从最新binlog开始,跳过历史数据)和 timestamp(从指定时间点开始)等选项,适用于不同场景。
  • sink.buffer-flush.*:这两个参数控制了写入Doris的批次行为。增大max-rowsinterval能提升吞吐,但会增加端到端延迟,需要根据业务容忍度权衡。
  • Doris Connector 的 sink.enable-delete:默认是 true,这确保了CDC捕获到的DELETE操作能在Doris中正确删除数据(通过发送DELETE SIGN)。如果你的场景只有新增和更新,可以关闭以提升少许性能。

3. 应对生产级挑战:分库分表与Schema变更

真实的业务场景远比单表同步复杂。接下来,我们探讨两个最常见的生产级问题及其解决方案。

挑战一:分库分表合并 互联网业务中,用户表 user_info 很可能被水平拆分成 user_info_001user_info_100 等多个分片。在数仓层,我们通常希望将它们合并成一张宽表。Flink CDC 可以优雅地解决这个问题。

-- 在Flink SQL中,我们可以使用UNION ALL来合并多个CDC源表
CREATE TABLE mysql_user_shard_source (
    `id` BIGINT,
    `name` VARCHAR(255),
    `email` VARCHAR(255),
    `create_time` TIMESTAMP(3),
    `shard_id` INT METADATA FROM 'table_name' VIRTUAL -- 元数据列,获取来源表名
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'mysql',
    'port' = '3306',
    'username' = 'root',
    'password' = '123456',
    'database-name' = 'user_db',
    'table-name' = 'user_info_.*', -- 使用正则表达式匹配多个表!
    'scan.startup.mode' = 'initial',
    'server-time-zone' = 'Asia/Shanghai'
);

-- 写入Doris时,可以保留shard_id用于溯源,也可以直接合并
INSERT INTO doris_user_sink
SELECT
    id AS user_id,
    name AS username,
    email,
    CAST(create_time AS TIMESTAMP) AS created_at,
    NOW() AS updated_at, -- 合并后统一更新时间
    0 AS is_deleted
FROM mysql_user_shard_source;

通过 table-name 参数使用正则表达式 user_info_.*,Flink CDC会自动监听所有匹配的表。METADATA 关键字让我们能获取到来源表的物理名称,便于数据处理或打标签。

挑战二:上游Schema变更 业务加字段、改字段类型是常态。传统的基于Canal+Kafka的方案,处理Schema变更非常麻烦。Flink CDC 2.3+ 版本对此提供了很好的支持。

当MySQL源表执行 ALTER TABLE ADD COLUMN age INT 后,Flink CDC作业默认会失败,因为Source的Schema与之前注册的已不匹配。为了自动化处理,我们需要在作业中启用Schema Evolution功能。

CREATE TABLE mysql_user_source_evolution (
    -- ... 字段定义 ...
) WITH (
    'connector' = 'mysql-cdc',
    -- ... 其他参数 ...
    'scan.incremental.snapshot.chunk.key-column' = 'id', -- 确保全量阶段有唯一键
    'debezium.schema.history.internal' = 'io.debezium.relational.history.MemorySchemaHistory', -- 内存存储Schema历史(仅测试)
    -- 生产环境建议使用Kafka或File存储Schema历史,如:
    -- 'debezium.schema.history.internal.kafka.bootstrap.servers' = 'kafka:9092',
    -- 'debezium.schema.history.internal.kafka.topic' = 'schema-changes.user_db'
);

在作业重启时,Flink CDC会从Schema历史存储中读取最新的表结构,自动更新Source的Schema定义。但是,下游Doris表的Schema也需要同步变更。目前,这需要手动或通过外部调度工具在Doris中执行相应的 ALTER TABLE 语句。一个实用的做法是,监控Flink作业的失败事件,触发一个自动化脚本去修改Doris表结构,然后重启Flink作业。社区也有一些基于Flink CDC Event的自动化方案在探索中。

提示:对于频繁变更的Schema,可以在Doris设计表时,将可能新增的维度字段设计为Map<String, String>类型的动态列,或者使用Doris 2.1+的Variant数据类型来存储半结构化JSON,上游Schema变更时只需修改写入的JSON内容即可,灵活性更高。

4. 数仓分层与实时数据建模

数据同步到ODS层(原始数据层)只是第一步。一个健壮的数仓需要清晰的分层。利用Flink强大的流式SQL计算能力和Doris的物化视图,我们可以在实时流上构建轻度汇总的DWD(明细数据层)甚至DWS(汇总数据层)。

方案一:Flink流上聚合后写入Doris 例如,我们需要实时统计每分钟的新增用户数。

-- 在Flink SQL中,基于CDC源表进行窗口聚合
CREATE TABLE doris_user_per_min_sink (
    `window_start` TIMESTAMP,
    `user_count` BIGINT
) WITH (
    'connector' = 'doris',
    'fenodes' = 'doris-fe:8030',
    'table.identifier' = 'doris_db.agg_user_per_min',
    'username' = 'root',
    'password' = ''
);

INSERT INTO doris_user_per_min_sink
SELECT
    TUMBLE_START(create_time, INTERVAL '1' MINUTE) AS window_start,
    COUNT(*) AS user_count
FROM mysql_user_source
WHERE deleted = FALSE -- 过滤已删除用户
GROUP BY TUMBLE(create_time, INTERVAL '1' MINUTE);

方案二:Doris物化视图(Materialized View) 对于更复杂的多表关联或历史全量数据的聚合,在Doris中创建物化视图是更轻量级的选择。物化视图会在底层数据变更时自动更新。

-- 在Doris中创建一张订单明细表(假设已通过CDC同步)
CREATE TABLE order_detail (
    `order_id` BIGINT,
    `user_id` BIGINT,
    `amount` DECIMAL(10,2),
    `order_time` DATETIME
) UNIQUE KEY(`order_id`) DISTRIBUTED BY HASH(`order_id`) BUCKETS 10;

-- 创建一张用户维度表
CREATE TABLE dim_user (
    `user_id` BIGINT,
    `user_name` VARCHAR(50),
    `city` VARCHAR(50)
) UNIQUE KEY(`user_id`) DISTRIBUTED BY HASH(`user_id`) BUCKETS 10;

-- 为“用户城市维度订单金额汇总”创建物化视图
CREATE MATERIALIZED VIEW mv_city_order_agg
AS
SELECT
    u.city,
    DATE(o.order_time) AS order_date,
    COUNT(o.order_id) AS order_cnt,
    SUM(o.amount) AS total_amount
FROM order_detail o JOIN dim_user u ON o.user_id = u.user_id
GROUP BY u.city, DATE(o.order_time);

创建物化视图后,查询 SELECT city, order_date, total_amount FROM mv_city_order_agg WHERE ... 时,Doris的优化器会自动路由到物化视图上,直接查询预计算好的结果,速度极快。这种方式将计算压力从Flink转移到了Doris,适合维度稳定、聚合逻辑明确的场景。

两种方案对比:

特性Flink流上聚合Doris物化视图
计算位置流计算引擎(Flink)分析型数据库(Doris)
计算触发事件驱动,延迟极低(秒级内)数据变更驱动,有一定延迟(通常秒级)
灵活性高,可进行复杂流式处理(如会话窗口、CEP)中,限于SQL定义的聚合和Join
维护成本需维护Flink作业状态Doris自动维护,无需额外作业
适用场景对延迟极度敏感、逻辑复杂的实时指标维度稳定的汇总报表、加速固定模式查询

在实际项目中,我通常采用混合模式:对核心的、延迟要求极高的实时风控或监控指标,用Flink流上计算;对面向分析师的报表和即席查询,用Doris物化视图来加速。这样既保证了核心链路的性能,又降低了整体架构的复杂度。

5. 性能调优、监控与故障处理

一个生产系统,除了能跑起来,更要跑得稳、跑得快。这里分享几个关键的调优和运维点。

Doris侧调优

  1. 分区分桶策略:这是影响Doris性能的基石。对于时间增长的表(如订单),一定要按时间分区(PARTITION BY RANGE)。分桶的列应选择查询中高频过滤或Join的列,桶数量建议是BE节点数的10-20倍,保证数据均匀分布。
    CREATE TABLE order_detail_optimized (
        `order_id` BIGINT,
        `user_id` BIGINT,
        `amount` DECIMAL(10,2),
        `order_time` DATETIME
    )
    UNIQUE KEY(`order_id`)
    PARTITION BY RANGE(`order_time`) () -- 先创建空分区,后续动态增加
    DISTRIBUTED BY HASH(`user_id`) BUCKETS 20
    PROPERTIES (
        "dynamic_partition.enable" = "true", -- 启用动态分区
        "dynamic_partition.time_unit" = "DAY",
        "dynamic_partition.start" = "-7", -- 保留最近7天
        "dynamic_partition.end" = "3",
        "dynamic_partition.prefix" = "p",
        "replication_num" = "3"
    );
    
  2. 索引加速:除了主键索引,Doris的倒排索引(Inverted Index) 对文本字段的模糊查询(LIKE)和等值过滤有奇效。Bloom Filter索引 对高基数列的等值过滤裁剪效果显著。在建表时或后期通过 ALTER TABLE 添加。
    ALTER TABLE user_info_sink ADD INDEX idx_email (`email`) USING INVERTED;
    ALTER TABLE order_detail ADD INDEX idx_amount (`amount`) USING BLOOM_FILTER;
    

Flink CDC调优

  1. 并行度:CDC Source的并行度决定了读取历史数据的速度。对于分库分表,可以设置并行度等于分表数,实现并发读取。增量阶段,并行度受限于MySQL binlog reader的数量,通常不需要太高。
  2. Checkpoint间隔:这是故障恢复的“存档点”。间隔太短(如10秒)会给State Backend(如RocksDB)带来压力;间隔太长(如10分钟)则恢复时可能重放大量数据。根据业务容忍度和系统负载,设置在1-5分钟是比较常见的。
  3. 资源分配:确保TaskManager有足够的堆内存,特别是当同步表数量多、历史数据量大时。给RocksDB State Backend分配足够的本地SSD空间。

监控与故障处理

  • Flink监控:通过Web UI或对接Prometheus,密切关注 numRecordsInPerSecond(输入速率)、numRecordsOutPerSecond(输出速率)、lastCheckpointDuration(Checkpoint耗时)以及 pendingRecords(积压记录数)。积压数持续增长,通常意味着下游写入(Doris)成为瓶颈。
  • Doris监控:使用Doris内置的Web UI或 SHOW PROC 命令,监控BE的磁盘使用率、Tablet版本数、Compaction状态。频繁的Schema变更或大量小批量写入可能导致Compaction跟不上,产生大量小文件,影响查询性能。可以通过调整 cumulative_compaction_num_threads_per_disk 等BE参数来优化。
  • 典型故障
    • Doris写入报错“Tablet writer write failed”:检查BE日志,常见原因是单次写入批次太大(调整Flink Connector的 sink.buffer-flush.max-rows)或BE磁盘满了。
    • CDC源表无法捕获删除:检查MySQL binlog格式是否为 ROW,以及用户是否有 REPLICATION SLAVE, REPLICATION CLIENT 权限。
    • 数据延迟增大:首先检查网络和下游Doris集群负载。其次,检查Flink作业是否有反压(Backpressure)。如果是Doris Compaction慢,可以临时增加BE资源或调整Compaction参数。

经过以上几个环节的打磨,一个具备生产可用性的实时数仓管道就基本成型了。从我的经验来看,Flink+Doris的组合真正做到了“1+1>2”。它用相对简单的架构,解决了实时数据同步、存储、计算和查询的全链路问题,让开发团队能更专注于业务逻辑本身,而不是日夜不停地“救火”和“缝补”复杂的系统。

更多推荐