十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

一条 SQL 顶一条 Flink 链路?Doris Streaming Job 持续导入全景解析

一条 SQL 顶一条 Flink 链路?Doris Streaming Job 持续导入全景解析 一、先从一个老问题说起把业务库的数据实时同步进 Doris 做分析传统链路长这样MySQL → Canal/Debezium → Kafka → Flink → Doris这条链路上的每一个环节都很成熟但问题也很现实组件多Kafka 集群、Flink 集群、Connector每一层都要部署、监控、值班延迟叠加数据每多流经一个系统端到端延迟就叠加一段语义难对齐exactly-once 要靠 Flink 两阶段提交 Sink 端配合出问题排查链路长。很多场景的诉求其实很简单——“把上游的变更持续搬进 Doris”并不需要复杂的流式计算。为这个诉求维护一整套流处理基础设施性价比很低。Doris 4.1 给出的答案是Streaming Job把 CDC 读取能力直接内置进 Doris 内核省掉所有中间组件创建同步任务只需要一条 SQL。二、Streaming Job 是什么Streaming Job 是 Doris 内置的持续导入作业提交之后Doris 会持续运行这个作业实时读取上游数据源的增量数据并写入 Doris 表。它复用了 Doris 已有的 Job 框架语法上只是在CREATE JOB里加一个ON STREAMINGCREATEJOB my_jobONSTREAMINGDOINSERTINTOdb1.tbl1SELECT*FROMcdc_stream(...);目前支持三种数据源官方文档持续导入概览数据源支持版本单表同步整库同步MySQL5.6、5.7、8.0.x✓✓PostgreSQL14、15、16、17✓✓S3-✓-版本边界需要说清楚MySQL / PostgreSQL 的 CDC 持续同步能力自 4.1.0 起完整支持S3 持续导入更早一些在 4.0.x 已经落地4.1.0 Release Notes。三、两种同步模式先搞清楚这个再动手Streaming Job 有两种实现机制完全不同的同步模式注意这不是同步一张表还是多张表的区别——自动建表同步也可以通过include_tables只同步一张表。能力维度SQL 映射同步自动建表同步底层机制Job TVFINSERT INTO tbl SELECT * FROM cdc_stream()Job 原生 DDLFROM MYSQL (...) TO DATABASE db目标一张已存在的 Doris 表一个 Doris database 容器建表方式需预建首次同步自动创建主键表数据加工支持列映射、过滤、转换SELECT 的完整表达能力原样复制不支持 ETL语义保证exactly-onceat-least-once靠主键表幂等兜底权限LoadLoad Create自动建表时适用场景需要列裁剪、字段重命名、类型转换、条件过滤整库/多表镜像复制表结构自动跟随上游选型原则很简单同步过程中要对数据做加工或者对精确一次语义有硬性要求 →SQL 映射同步希望一条 SQL 把一组表或整库镜像过来下游表自动建、结构自动跟随 →自动建表同步数据源是 S3 → 只有 SQL 映射同步S3 TVF一条路。四、上手实战4.1 Demo 1MySQL 整库同步一条 SQL前置条件MySQL 开启 binlog、账号有读 binlog 权限、JDBC 驱动 jar 可引用文件名 / 本地路径 / HTTP 地址均可。CREATEJOB multi_table_syncONSTREAMINGFROMMYSQL(jdbc_urljdbc:mysql://127.0.0.1:3306,driver_urlmysql-connector-java-8.0.25.jar,driver_classcom.mysql.cj.jdbc.Driver,userroot,password123456,databasetest,include_tablesuser_info,order_info,offsetinitial)TODATABASEtarget_test_db(table.create.properties.replication_num1-- 单 BE 环境需要设为 1);几个关键点include_tables逗号分隔留空即同步整库offset initial表示先全量初始化再自动切换增量只想要增量就写latestTO DATABASE后的table.create.properties.*控制自动建表的表属性。创建之后用jobs()表函数观察运行状态select*fromjobs(typeinsert)whereExecuteTypeSTREAMING\G输出里最值得关注的三个字段示例来自官方文档Status: RUNNING CurrentOffset: {ts_sec:1765284495,file:binlog.000002,pos:9350, ...} LoadStatistic: {scannedRows:24,loadBytes:1146,fileNumber:0,fileSize:0}CurrentOffset直接展示当前消费到的 binlog 位点文件 position同步进度一目了然——排障时先看它再看ErrorMsg。4.2 Demo 2S3 持续导入S3 场景是目录监控模式Doris 持续探测指定路径下新增的文件并自动导入。CREATEJOB s3_jobONSTREAMINGDOINSERTINTOdb1.tbl1SELECT*FROMS3(uris3://bucket/demo/*.csv,formatcsv,column_separator,,s3.endpointhttps://s3.ap-southeast-1.amazonaws.com,s3.regionap-southeast-1,s3.access_key...,s3.secret_key...);S3 模式有两个攒批参数通过 Job 的PROPERTIES设置任一条件满足即触发一次写入参数默认值说明s3.max_batch_files256累计文件数达到该值触发写入s3.max_batch_bytes10GB累计字节数达到该值触发写入可配范围 100MB ~ 10GBmax_interval10秒上游没有新数据时的空闲调度间隔文件全部导入完成后作业会进入FINISHED状态。五、原理浅析它在内核里是怎么跑的5.1 整体架构FE 是调度中枢Job 的元信息、状态机、offset 都持久化在 FE。StreamingJobSchedulerTask负责按状态机驱动作业流转StreamingInsertTask是实际干活的最小单元BE 侧跑 CDC ReaderCDC 场景集成了 Flink CDC 的读取能力全量 snapshot 增量 binlog/WAL读到的数据经Stream Load写入 Doris增量阶段 Reader 会绑定固定 BE源码里每个 Job 维护一个boundBackendId增量阶段优先在同一台 BE 上复用 CDC Reader避免重复初始化BE 变更时会重新绑定并持久化。5.2 调度时间驱动 事件驱动双引擎Job 调度时间驱动复用 Job 框架的时间轮定期产生调度子任务任务调度事件驱动由上一个 Task 完成的回调驱动下一次调度数据来得快时不会被固定间隔卡住。上游没数据时也不是空转干等——按max_interval默认 10 秒的间隔空闲轮询。5.3 exactly-once 是怎么实现的核心是持久化 offset 两阶段任务验证offset 滞后提交只有当数据在 Doris 里可见且持久化之后offset 才会提交。FE 侧通过StreamingTaskTxnCommitAttachment把 offset 随事务一起落盘——offset 推进和数据可见性是同一个事务的副产品单调任务 ID 防重放每个 Task 携带单调递增 ID调度器拒绝任何重复或乱序的 Task从机制上消除重放风险。反过来说这也解释了为什么自动建表同步只能做到 at-least-once整库镜像模式下 CdcClient 直接消费上游变更流中断恢复后可能重放一小段数据——但目标表是主键表重放的数据会被幂等覆盖最终效果等价于 exactly-once。5.4 几个源码里才能看到的细节autoResume 的退避策略作业因故障进入PAUSED后调度器按指数退避自动恢复第 n 次重试等待2^n × 10秒封顶 300 秒重试超过 5 次后固定 300 秒一轮。临时网络抖动基本都能自愈自动恢复有预算上限streaming_job_max_auto_resume_count默认 10可动态修改。重试耗尽后失败原因会被改写为CANNOT_RESUME_ERR必须人工RESUME JOB介入——这是防止无限重试打爆上游的保护阀。另外手动暂停和失败行数超限max_filter_ratio不会触发自动恢复Task 超时放宽Streaming Task 的insert_timeout和query_timeout默认被放宽到 30 分钟避免长快照阶段的任务被常规超时误杀FE 侧 split 缓冲有软上限已生成未消费的 split 最多缓冲 512 个MAX_PENDING_SPLITS防止 FE 内存被积压数据打满。六、运维与可观测性6.1 状态机PENDING等待调度→ RUNNING执行中→ FINISHED源消费完毕如 S3 文件导完 ↘ PAUSED子任务失败自动暂停或人工暂停PAUSED 之后两条路autoResume 按退避策略自动拉起回到 PENDING或者人工排查后RESUME JOB。6.2 日常运维命令-- 查看所有 Streaming 作业重点看 Status / CurrentOffset / LoadStatistic / ErrorMsgselect*fromjobs(typeinsert)whereExecuteTypeSTREAMING;-- 查看某个作业的所有子任务看 RunningOffset 和单 Task 的失败信息select*fromtasks(typeinsert)wherejobIdjob_id;-- 暂停手动暂停不会被 autoResume 唤醒PAUSE JOBWHEREjobnamejob_name;-- 恢复RESUME JOBWHEREjobnamejob_name;-- 修改如上游账号密码轮换后更新连接信息ALTERJOBFORjob_namePROPERTIES(...);-- 删除DROPJOBWHEREjobnamejob_name;6.3 FE 配置参数参数说明max_streaming_job_num最大 Streaming 作业数默认 1024job_streaming_task_exec_thread_num执行 StreamingTask 的线程数max_streaming_task_show_count内存中每个 Job 最多保留的 Task 记录数默认 100streaming_job_max_auto_resume_count自动恢复重试预算默认 10可动态调整七、限制与坑仅支持带主键的上游表目标表必须是 Unique Key 模型——自动建表模式会自动建成主键表SQL 映射模式需要你自己建对DDL 同步能力有限MySQL 上游的 DDL 目前不同步需要手工调整 Doris 表结构PostgreSQL 自 4.1 起支持同步ADD COLUMN/DROP COLUMN但列类型变更、RENAME COLUMN、约束/索引/分区变更都不会同步语义差异要记牢SQL 映射同步是 exactly-once自动建表同步是 at-least-once靠主键幂等对外宣传时不要笼统说精确一次MySQL 连接常见报错Public Key Retrieval is not allowed——JDBC URL 加allowPublicKeyRetrievaltrue或把用户认证方式改为mysql_native_password八、结尾Streaming Job 的意义不在于又多了导数方式而在于 Doris 把数据接入这件事的默认形态从搭一条外部管道变成了发一条 SQL。对于大量只需要镜像同步 轻量加工的实时数仓场景架构里可以直接少掉 Kafka 和 Flink 两层。这个方向还在快速演进社区正在讨论的 Table Stream MTMV 增量计算#65418——如果落地Doris 里的流就不只是导入而是完整的增量计算链路。从流式导入到流式数仓这条路才刚开始。
返回列表