ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MySQL数据迁移Doris实战:Stream Load与Flink CDC一次性讲透

MySQL数据迁移Doris实战:Stream Load与Flink CDC一次性讲透 上个月帮一个朋友把跑了快三年的 MySQL 订单库迁到 Doris单表 3000 万行、5.6GB本以为就是导出再导入结果被各种细节问题整整折腾了三天。这篇文章不是官方文档的复述更像是一份真实迁移记录把 MySQL 数据快速导入 Doris 过程中必须搞清楚的技术点一次说透。如果你正在纠结 Doris 和 ClickHouse 怎么选或者已经决定用 Doris、但不知道用哪种方式导数据建议先把这个流程消化一遍。先说结论一次性全量导入我最推荐 Stream Load持续增量同步最推荐 Flink CDC如果只是临时拉几张百万级的小表验证外部表 INSERT INTO SELECT 就够了。下面把每个选择都拆开讲清楚。1. 迁移前先想清楚Doris 到底解决了什么问题1.1 MySQL 在大数据量分析下的真实瓶颈MySQL 是行式存储加聚簇索引点查单行、小范围扫描非常快但一旦涉及全表扫描、多表关联、高频聚合性能会随着数据量上涨迅速崩掉。千万行可能还能忍过亿之后一条 group by 跑几十秒甚至几分钟都很正常。加索引、做主从复制、上中间层缓存都只是缓解读压力并没有改变 MySQL 不适合大规模并行扫描的底层事实。我见过一个团队把 5 亿行订单留在 MySQL 里做经营报表每个查询都要等一分钟以上业务方天天催数。后来引入 Doris 之后同样的查询从几十秒压到了秒级。这个提升不是来自数据库“更聪明”而是列式存储加 MPP 并行扫描带来的结构性优势。Doris 会把一张表按分桶拆成多个 tablet 分布在多台 BE 上查询时多节点同时扫描效果自然不一样。1.2 Doris 和 ClickHouse 怎么选很多人会拿 Doris 和 ClickHouse 对比。我的经验是如果源端就是 MySQL并且业务有主键更新、删除、点查这类需求Doris 的适配成本通常比 ClickHouse 低很多。ClickHouse 的 MergeTree 在事件流、大宽表、超大规模日志分析上很强但它的更新和删除能力非常弱修改主键更是成本极高。Doris 的 Unique Key 模型天然支持 UPSERT 和 DELETE建表方式、查询语法也更接近 “分布式 MySQL” 的感觉工具链上甚至可以直接用 Navicat 连 Doris 的 9030 端口。反过来如果你只是要一个列式存储跑海量日志分析没有高频更新ClickHouse 也完全够用。选型的核心不是比谁更“先进”而是看哪边的语义贴近你的业务。我之前见过有人非要用 ClickHouse 承载订单数据结果每次状态流转都要 insert 一条新版本查询时还要搞 argMax 这类操作维护成本非常高。换成 Doris 的 Unique Key 表一个 INSERT OVERWRITE 就能解决流程瞬间清爽。1.3 先明确场景边界再决定要不要搬适合用 Doris 承接 MySQL 数据的场景基本是这几类离线报表、数仓 ODS/DWD 层、用户行为分析、需要同时支持明细查询和聚合计算的大表。不适合的也很明确高频事务写入、强一致点查、数据量不到百万级的小规模业务。小表留在 MySQL 里毫无问题没必要为了“用 Doris”而引入一套新技术栈。建表模型一旦选错后面返工很痛。所以我的习惯是先写清楚业务查询的模式是需要保留全部明细还是按主键覆盖还是直接做预聚合。这个决定直接对应 Doris 的 Duplicate Key、Unique Key、Aggregate Key 三种模型我会在下一章展开。2. 准备阶段MySQL 开 binlog、Doris 最小集群搭起来2.1 MySQL 端这 4 行配置别省MySQL 怎么搭建我不展开系统包管理器或者 Docker 都可以但有两个点必须确认第一个是版本5.7 和 8.0 都能和 Doris 对接第二个是 binlog 配置如果你后面想用 Flink CDC 做增量同步binlog 没开就等于没有数据源。在 my.cnf 里加上这几行log_binON server_id1 binlog_formatROW binlog_row_imageFULL expire_logs_days7binlog_format 必须改成 ROW。原因很简单Doris 侧的各种同步工具要解析 binlog 拿到每一行变更前后的完整数据STATEMENT 或 MIXED 格式只记录 SQL 语句或混合内容很多字段值根本还原不出来。binlog_row_imageFULL 是保证即使只更新了一个字段binlog 里也包含整行所有列的值。expire_logs_days 设置成 7 是为了避免磁盘被撑爆也够同步任务追赶进度了。再创建一个专门用于同步的账号CREATE USER sync% IDENTIFIED BY sync_pass; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO sync%; FLUSH PRIVILEGES;全量导出只需要 SELECT 权限但 CDC 增量同步必须要有 REPLICATION 相关权限所以建议一开始就按这个口径给。如果你是用 Docker 装的 MySQL记得把 3306 端口映射出来并且确认容器重启后 binlog 配置仍然生效。2.2 Doris 端FE 加 BE 两个角色跑起最小集群Doris 集群的最小单元是一个 FE 加一个 BE。FE 负责接收 SQL、管理元数据BE 负责存储数据和计算。下载 Apache Doris 二进制包后目录里会有 fe 和 be 两个子目录分别修改配置即可。FE 的核心配置是 fe/conf/fe.conf里面通常需要设置 priority_networks尤其是机器有多个网卡的时候。不设置的话BE 注册时可能拿到错误的 IP导致节点状态一直是 false。启动命令是sh bin/start_fe.sh --daemon然后用 MySQL 客户端连 FE 的查询端口mysql -h 127.0.0.1 -P 9030 -u rootDoris 初始 root 用户没有密码直接用空密码登录。接着启动 BEsh bin/start_be.sh --daemon回到 FE 的 MySQL 命令行里挂载 BEALTER SYSTEM ADD BACKEND 127.0.0.1:9050; SHOW BACKENDS;如果 Alive 这一列是 true就说明 BE 正常接入。常见的坑集中在端口没放行、BE 数据目录没有写权限以及磁盘容量不足这三个点上。测试环境可以把 be.conf 里的 mem_limit 调小默认是机器内存的 90%如果这台机器还跑着 MySQL很容易内存超限。我一般设成 60% 左右。Doris 常用端口对应关系组件端口用途FE8030HTTP 服务Stream Load 提交入口FE9030MySQL 协议端口SQL 查询连接BE8040BE 的 HTTP 端口Stream Load 实际接收端BE9050BE 心跳端口负责向 FE 注册BE9060BE Thrift 端口用于执行查询2.3 建表模型不是随便建一张表就行Doris 建表前必须先理解三种模型。Duplicate Key 保留所有导入的明细行适合日志、事件这类没有更新语义的数据。Unique Key 按主键去重新数据覆盖旧数据适合从 MySQL 迁过来的订单、用户等有主键的表。Aggregate Key 则是在相同 key 上做预聚合适合指标统计比如 SUM、MAX、MIN、REPLACE。迁移订单表我一般用 Unique Key 模型CREATE TABLE dwd_order ( order_id BIGINT NOT NULL COMMENT 订单id, user_id BIGINT NOT NULL COMMENT 用户id, order_time DATETIME NOT NULL COMMENT 下单时间, amount DECIMAL(10, 2) NOT NULL COMMENT 订单金额, status VARCHAR(20) NULL COMMENT 订单状态 ) UNIQUE KEY(order_id) DISTRIBUTED BY HASH(order_id) BUCKETS 16 PROPERTIES ( replication_num 1 );分桶数 BUCKETS 我这里写了 16测试环境完全够用。经验值是按单 bucket 数据量控制在 1GB 到 2GB 来估算并且最好是 2 的幂。生产环境副本数建议设成 2 或 3我这里 replication_num 设 1 是为了先在单机验证流程。Unique Key 表的列顺序最好让主键在前分桶键也尽量选均匀分布的列比如 order_id避免 hash 倾斜导致某个 BE 压力特别大。MySQL 字段到 Doris 字段的类型映射也要提前规划。常用映射关系是这样的MySQL 类型Doris 类型说明intINT直接对应bigintBIGINT注意无符号 bigint 会溢出建议改用 DECIMAL(20,0)varchar(n)VARCHAR(n) 或 STRING长度要按实际最大值预留datetime / timestampDATETIME如果 MySQL 有微秒精度Doris 建表用 DATETIME(6)dateDATE直接对应decimal(p,s)DECIMAL(p,s)精度保持一致text / longtextSTRINGDoris 里 STRING 可以存大字段3. 导入方式对比外部表、DataX、Stream Load 和 Flink CDC 到底用哪个3.1 每种方式的核心定位先说外部表。Doris 可以直接建一张 MySQL 外部表然后用一条 INSERT INTO SELECT 把数据拉进来。这是最轻量的方案适合数据量不大、一次性验证的场景。我在本机做过测试一百万行以内的表用这种方式几分钟就能导完。但数据量超过一定规模后所有数据都要经过 FE 节点中转很容易成为瓶颈所以不建议拿它做大宗迁移。DataX 是阿里巴巴的开源数据同步工具通过 mysqlreader 读 MySQL再通过 doriswriter 写 Doris。如果你需要定时大批量全量同步DataX 是稳定选择它天然支持分片、限流、断点重跑适合接调度平台。缺点是 doriswriter 需要额外准备而且不同 DataX 版本对不同 Doris 版本的兼容性需要踩一轮坑。社区里常常有人卡在 writer 插件和 Doris 版本不匹配的问题上所以用之前一定要确认好版本组合。Stream Load 是 Doris 官方最推荐的批量导入方式。它基于 HTTP 协议客户端把数据文件直接发给 FEFE 再重定向到 BE 实际写入。Stream Load 速度快、操作简单几十 GB 到 TB 级的数据都能处理也是我这次迁移的主力方案。Flink CDC 是增量同步的默认解。Flink CDC 连接 MySQL 的 binlog把变更记录解析成 changelog 流再通过 Doris Connector 写入表。Doris 的 Unique Key 模型能直接消费这种 upsert/delete 语义搭建好之后可以实现秒级延迟的准实时同步。3.2 一张表看懂取舍导入方式适合数据量实时性复杂度典型场景外部表 INSERT百万级一次性低临时验证、小表迁移DataX百GB/TB级定时中每日全量、可断点同步Stream LoadGB/TB级批次低全量迁移、批量增量Flink CDC持续小批量秒级准实时较高实时数仓、整库同步3.3 我给出的默认组合如果让我接手一个全新迁移我的默认方案是全量用 Stream Load增量用 Flink CDC外部表只用来做流程验证和临时抽数DataX 只在已有团队强依赖它的情况下才引入。原因是 Stream Load 足够快且操作非常直接Flink CDC 则把增量同步这条最麻烦的链路标准化了。准备阶段就在 MySQL 打开 binlog后面切增量时不用重新配置环境可以省掉很多来回折腾的时间。4. Stream Load 批量导入命令、参数和一次实测4.1 先把 MySQL 数据导成标准 CSVStream Load 并不能直接读 MySQL你需要先把数据导出成文件。最不推荐的做法是直接mysql -e select * from table data.csv因为 MySQL 客户端的文本输出对换行符、制表符、NULL 值的处理都不可控导出来后很容易出现列错位和行数不一致的问题。我更推荐用一个小脚本从 MySQL 读出数据后按标准 CSV 写到本地。比如用 Pythonimport csv import pymysql conn pymysql.connect( host127.0.0.1, port3306, usersync, passwordsync_pass, databasemydb, charsetutf8mb4, cursorclasspymysql.cursors.SSCursor ) with conn.cursor() as cur: cur.execute(SELECT order_id, user_id, order_time, amount, status FROM orders) with open(orders.csv, w, newline, encodingutf-8) as f: writer csv.writer(f, quotingcsv.QUOTE_MINIMAL, lineterminator\n) for row in cur: writer.writerow(row)用 SSCursor 可以避免一次性把 3000 万行全加载进内存。导出时如果字段里包含逗号、引号、换行符csv 模块会自动做引号转义这份文件才是 Stream Load 能可靠解析的输入。如果你暂时没有 Python 环境也可以用 MySQL 的SELECT ... INTO OUTFILE在服务端直接生成 CSV但要提前确认 secure_file_priv 权限并且导出的分隔符要和你后面指定的 column_separator 保持一致。4.2 提交 Stream Load 的完整姿势拿到 CSV 之后用 curl 提交 Stream Load。这里有一个非常容易踩的坑登录认证必须配合--location-trusted参数。因为请求先发到 FE 的 8030 端口FE 会返回一个重定向到具体 BE 的 8040 端口不加--location-trusted的话重定向时会丢掉认证信息得到一堆 401 或者循环重定向很多人第一次跑都会被这个卡住。正确的命令格式curl --location-trusted \ -u root: \ -H Expect: 100-continue \ -H format: csv \ -H column_separator: , \ -H columns: order_id,user_id,order_time,amount,status \ -T /tmp/orders.csv \ http://127.0.0.1:8030/api/dwd/dwd_order/_stream_load逐个说明URL 固定是http://FE_HOST:8030/api/数据库名/表名/_stream_load-u root:是登录凭据root 用户密码为空时冒号后什么都不写format: csv指定文件格式column_separator: ,和 CSV 文件里的分隔符保持一致columns把 CSV 里的列名映射到 Doris 表字段顺序必须一致多出的列还可以通过这个参数做过滤或函数转换。4.3 返回结果、label 和 max_filter_ratio执行成功后Stream Load 会返回一段 JSON类似这样{ TxnId: 1001, Label: orders_20240601_001, Status: Success, NumberTotalRows: 30000000, NumberLoadedRows: 30000000, NumberFilteredRows: 0, LoadBytes: 5632000000, LoadTimeMs: 720000, ErrorURL: }我建议每次提交时手动指定一个带时间语义的 label比如orders_20240601_001。Doris 会把 label 当作导入事务的唯一标识同一个 label 重复提交时会直接返回 Duplicate label 错误而且这个错误不是一个 bug而是 Doris 的幂等保护机制防止你重复提交同一个任务。max_filter_ratio是另一个关键参数默认是 0表示有任何一行解析失败整个任务都会失败。如果你的数据质量一般可以设置为 0.1 之类的值允许 10% 的脏数据被过滤掉而不影响整体任务。但我实际使用中建议谨慎调大因为高容错率会把类型转换错误、列错位这类问题悄悄掩盖掉等你核对数据时才发现少了一堆行排查反而更痛苦。先确保数据干净再考虑放宽这个参数。merge_type参数也很重要。Stream Load 针对 Unique Key 表默认是 MERGE 语义也就是按主键合并新数据覆盖旧数据。如果只是往 Duplicate Key 表追加明细可以用 APPEND。如果你要实现按主键删除某些行那需要在columns里声明删除列并配合 DELETE 语义来用。这个参数建议按实际情况设置别全都依赖默认值。4.4 性能预期与大文件并发思路我这次迁移的单表是 3000 万行、5.6GB本地 SSD 环境下单个 1.2GB 的 CSV 用 Stream Load 导入大约 45 秒折合每秒 20 到 25MB。同时开 5 个并发任务每个文件独立 label总耗时能压到 12 分钟左右。如果你手里是大文件可以先切分再并发split -l 5000000 orders.csv part_ ls part_*切分会生成 part_aa、part_ab 等文件然后用脚本循环提交 Stream Load。注意并发数不是越大越好BE 内存和磁盘 IO 都有上限我一般控制在 5 到 10 个并发之间。文件大小方面单个文件 200MB 到 1GB 比较合适太小会制造大量小事务反而拖累整体效率。5. Flink CDC 增量同步从 binlog 到 Doris 的准实时链路5.1 为什么 CDC 适合增量场景Stream Load 虽然快但毕竟是批量操作不能做到秒级延迟。业务量上来之后Doris 里的报表数据需要跟着 MySQL 随时更新这就轮到 Flink CDC 上场。它的核心逻辑是连接 MySQL binlog持续监听每一行数据的 insert、update、delete 操作然后通过 Doris Connector 把这些变更写到 Doris。由于 Doris 的 Unique Key 表天然支持覆盖和删除MySQL 里的主键更新在 Doris 端就是直接覆盖对应主键逻辑上不需要额外转换。5.2 版本选型与 jar 包准备Flink CDC 是组合链路版本兼容最容易出问题。我实测过一套能跑通的组合Flink 1.18flink-sql-connector-mysql-cdc 2.4.0flink-doris-connector 1.6.0。三者都放在 Flink 的 lib 目录下或者放到 SQL Client 的 classpath 中。不同版本的 connector 之间偶尔会有类冲突或者某个方法签名对不上所以不要随便混搭版本最好先在测试环境验证一遍再上线。jar 包可以从 Maven 仓库直接下载。到 repo1.maven.org 上搜索 flink-sql-connector-mysql-cdc找到对应版本后再下载 flink-doris-connector。Doris 官方文档里会列出 connector 与 Doris 版本的对应关系比如 1.6.0 对应 Doris 2.x建议以官方给出的组合为准。5.3 Flink SQL 建表与同步任务准备好 jar 包后用 Flink SQL 定义源表和目标表。源表是 MySQL CDCCREATE TABLE orders_source ( order_id BIGINT, user_id BIGINT, order_time TIMESTAMP(3), amount DECIMAL(10, 2), status STRING, PRIMARY KEY(order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 127.0.0.1, port 3306, username sync, password sync_pass, database-name mydb, table-name orders, server-id 1001 );server-id必须和集群里其他同步任务保持唯一否则多个任务共享同一个 server_id 会导致从库断开或数据错乱。如果 MySQL 做了主从切换server_id 冲突会第一时间暴露出来所以这里要特别注意。目标表是 DorisCREATE TABLE orders_doris ( order_id BIGINT, user_id BIGINT, order_time TIMESTAMP(3), amount DECIMAL(10, 2), status STRING, PRIMARY KEY(order_id) NOT ENFORCED ) WITH ( connector doris, fenodes 127.0.0.1:8030, table.identifier dwd.dwd_order, username root, password , sink.label-prefix cdc_orders, sink.properties.format json, sink.properties.read_json_by_line true );sink.label-prefix是每次写入生成 Stream Load label 的前缀必须保证全局唯一多个同步任务用不同前缀避免冲突。sink.properties.format设置为 json 是为了让 Doris Connector 正确传递完整的 update/delete 语义配 CSV 格式也可以但 JSON 对 CDC 场景更直接。然后启动同步任务INSERT INTO orders_doris SELECT order_id, user_id, order_time, amount, status FROM orders_source;任务启动后Flink CDC 会先做一次一致性快照读取 MySQL 当前全量数据再平滑切到 binlog 增量。也就是说如果表数据量不是特别巨大你可以直接用 Flink CDC 完成一次“全量加增量”的同步不一定需要先单独跑 Stream Load。但数据量很大时我建议先用 Stream Load 导完历史数据在某个 binlog 位点附近启动 CDC然后补一小段增量即可这样任务启动压力更小。5.4 运维要点checkpoint、时区和 DDL 变更Flink CDC 任务必须配置 checkpoint否则任务重启后无法从 binlog 位点恢复会造成丢数据或者重复消费。checkpoint 间隔建议 30 秒到 60 秒。这个参数不在同步 SQL 里而是在 Flink 的配置文件中设置。时区也要提前对齐。MySQL 的 TIMESTAMP 和 Doris 的 DATETIME 对时区的处理逻辑不一样如果发现导入后时间差 8 小时优先检查两边的 time_zone 设置。Stream Load 里可以在请求头加timezone: 08:00Flink 任务则要保持运行环境的时区一致。还有一个常见问题是 DDL 变更。如果 MySQL 表新增了一个字段Doris 表和同步 SQL 没有同步更新任务会报字段缺失或解析失败。目前比较稳的方式是先在 Doris 端执行 ALTER TABLE 增加字段再重启 Flink CDC 任务。不要指望 CDC 自动做 schema evolution至少在成熟稳定之前手动处理会更安全。6. 导入后的核对与排错数据对不对、任务稳不稳6.1 怎么判断数据“没丢”导入完成的第一件事不是建报表而是核对数据。核对行数时要注意表模型的差异Duplicate Key 表里 count() 和 MySQL 一样是精确行数Unique Key 表因为主键覆盖最终行数可能小于 MySQL 源表行数这是正常的你应该对比的是count(distinct 主键)而不是 count()。我这次迁移的做法是同时核对行数和关键金额字段SELECT COUNT(*), SUM(amount) FROM dwd_order;在 MySQL 里跑同样的统计两边比对。这个方案够用但不算严格。更严格的校验可以分别算所有行的哈希比如SUM(CRC32(CONCAT_WS(|, order_id, amount, status)))两边结果一致说明字段值没有偏差。不过这种校验对字段顺序和 NULL 值敏感跑之前先确认两边类型转换规则一致尤其是 DATETIME 的精度Doris 默认到秒MySQL 如果存了微秒导入时会截断导致哈希值对不上。需要保留微秒就在建表时用 DATETIME(6)。6.2 一张表说清常见报错迁移过程中会碰到各种报错我把最典型的整理成了对照表报错或现象最常见原因解决办法Duplicate label同一个 label 重复提交每次任务换新 label不要复用Memory limit exceededBE 内存不足导入任务占用过高降低并发数检查 be.conf 的 mem_limitThe tablet write operation failedBE 掉线、磁盘满、副本数不足SHOW BACKENDS 检查节点状态清理磁盘空间Status: FailErrorURL 有内容部分行解析失败或类型不匹配打开 ErrorURL 定位具体行修数据或调整 max_filter_ratioFlink CDC 一直处于 init snapshot表太大快照阶段慢调大 Flink 内存或先用 Stream Load 导全量再接增量导入后时间差 8 小时时区不一致Stream Load 加 timezone 参数Flink 任务统一时区6.3 针对导入性能做后续调优第一批数据导完后如果发现查询性能不理想优先看三个点分桶数、副本数、Compaction 状况。分桶数太少会导致某些 tablet 数据量过大查询时单点压力明显副本数在测试环境设 1 很省资源但生产至少 2否则 BE 挂掉一块盘数据就丢了。Compaction 会自动把大量小文件合并成大文件一般不用手动干预但如果每天导入频繁可以观察 BE 日志里 compaction 线程是否有堆积。Stream Load 的并发参数也可以持续调整。数据量小且文件不多时2 到 3 个并发足够了数据量大时5 到 10 个并发是常见选择。如果 BE 内存吃紧优先降低并发而不是调大 max_filter_ratio因为过滤比例调大了容易掩盖真实数据问题。最后分享一个我自己的体会这套流程上线跑了两个多月出问题最多的时候其实不是导入本身而是两边字段类型没对齐和数据格式差异。所以别急着一次性把几十张表全迁完先挑一张有代表性的订单表把流程打通确认建表模型、全量导入、增量同步都稳定运行几天再批量复制。这个节奏虽然慢一点但后面的返工成本会小很多。
返回列表