跳到主要内容
最后 更新

持续导入概览

Doris 支持通过 Streaming Job 的方式,从多种数据源持续导入数据到 Doris 表中。提交 Job 后,Doris 会持续运行导入作业,实时读取数据源中的数据并写入到 Doris 表中。

提示

该功能自 4.1.0 版本起支持。

本文将帮助你解决以下问题:

  • 持续导入支持哪些数据源和同步模式?
  • SQL 映射同步与自动建表同步该如何选择?
  • 作业的运行状态如何流转、如何自动恢复?
  • 日常如何查看、暂停、恢复、删除导入作业?
  • 有哪些通用的 FE 与 Job 配置参数?

支持的数据源与同步模式

持续导入支持以下数据源和同步模式:

数据源支持版本SQL 映射同步自动建表同步配置指南
MySQL5.6、5.7、8.0.xMySQL CDC SQL 映射同步MySQL CDC 自动建表同步Amazon RDS MySQL · Amazon Aurora MySQL
PostgreSQL14、15、16、17PostgreSQL CDC SQL 映射同步PostgreSQL CDC 自动建表同步Amazon RDS PostgreSQL · Amazon Aurora PostgreSQL
OceanBaseMySQL 兼容模式-支持(自 4.1.4 版本起),语法见下方【OceanBase 数据源】-
S3-S3 持续导入--

上游列类型如何映射为 Doris 类型,见 MySQL / PostgreSQL 数据类型映射。OceanBase 复用 MySQL 的类型映射。

OceanBase 数据源

自 4.1.4 版本开始支持。

OceanBase 作为持续导入的 CDC 数据源,使用 FROM OCEANBASE (...) 子句,属性与 MySQL 数据源一致:

CREATE JOB oceanbase_sync ON STREAMING
FROM OCEANBASE (
"jdbc_url" = "jdbc:mysql://<host>:<port>",
"driver_url" = "<driver_jar_url>",
"driver_class" = "com.mysql.cj.jdbc.Driver",
"user" = "<user>",
"password" = "<password>",
"database" = "<ob_database>",
"include_tables" = "t1,t2",
"offset" = "initial"
)
TO DATABASE <doris_db> (
"table.create.properties.replication_num" = "1"
);

使用限制:

  • jdbc_url 必须以 jdbc:mysql:// 开头,否则报错 OceanBase jdbc_url must start with 'jdbc:mysql://'
  • 只支持 OceanBase 的 MySQL 兼容模式。作业创建时会执行 SHOW VARIABLES LIKE 'ob_compatibility_mode' 探测;Oracle 兼容模式会报错 OceanBase Oracle compatibility mode is not supported for streaming jobs
  • 不支持 schemaslot_namepublication_name 属性,指定后会报错 Property '<key>' is not supported for OceanBase
  • database 为必填项;
  • jdbc_url 的参数归一化规则与 MySQL 一致,详见 数据类型映射

如何选择同步方式

SQL 映射同步和自动建表同步是两种实现机制完全不同的持续导入方式,并非"表数量"的区别。自动建表同步也支持通过 include_tables 只同步一张表,因此选型应以能力需求为准。

能力对比

能力维度SQL 映射同步自动建表同步
底层机制Job + TVF(INSERT INTO tbl SELECT * FROM tvf()Job + 原生整库 DDL(FROM src TO DATABASE db
目标层级一张已存在的 Doris 表一个 Doris database 容器
同步范围单张表一张到多张到整库(由 include_tables 控制)
自动建表需预建首次同步自动创建主键表
SQL 灵活表达支持列映射、过滤、转换(SELECT 子句)原样复制,不支持 ETL
语义保证exactly-onceat-least-once
所需权限LoadLoad + Create(自动建表时)
典型适用场景需要列裁剪、字段重命名、类型转换、条件过滤的实时同步整库或一组表的镜像复制,希望下游表结构自动跟随上游

选型建议

  • 需要对数据做 SQL 加工,或对精确一次语义有严格要求 → 选 SQL 映射同步
  • 希望 Doris 自动建表、一次配置同步一组表 → 选 自动建表同步
  • 数据源是 S3 对象存储 → 只支持 SQL 映射同步(S3 TVF 方式)

作业状态流转

Streaming Job 在运行过程中会在以下状态之间迁移,SQL 映射同步和自动建表同步遵循同一套状态机

job-state-flow

状态说明

状态含义
PENDING作业已创建但尚未调度出子任务;等待下一次调度创建 StreamingTask
RUNNING已派生子任务并在执行中,从源端读取增量数据并写入 Doris
FINISHED源消费完成,作业终止。S3 TVF 文件全部导入完成后会进入该状态
PAUSED子任务执行失败,作业自动暂停并记录 failReason;可通过 select * from jobs(...)ErrorMsg 字段查看原因

自动恢复(autoResume)

作业进入 PAUSED 后,调度器会按指数退避策略定时尝试恢复,恢复时回到 PENDING 继续创建子任务。无需人工介入临时故障(网络抖动、上游短暂不可用等)会被自动消化。

不同场景下应使用的命令:

  • 立即恢复或排查故障后手动启动:使用 RESUME JOB
  • 彻底停止不再调度:使用 PAUSE JOB(手动暂停不会被 autoResume 唤醒)或 DROP JOB

通用操作

查看导入状态

查询所有 Streaming 类型的 Insert Job:

select * from jobs("type"="insert") where ExecuteType = "STREAMING";

返回结果列说明:

结果列说明
IDJob ID
NAMEJob 名称
DefinerJob 定义者
ExecuteTypeJob 调度的类型:ONE_TIME/RECURRING/STREAMING/MANUAL
RecurringStrategy循环策略。普通的 Insert 会用到,ExecuteType=Streaming 时为空
StatusJob 状态
ExecuteSqlJob 的 Insert SQL 语句
CreateTimeJob 创建时间
SucceedTaskCount成功任务数量
FailedTaskCount失败任务数量
CanceledTaskCount取消任务数量
CommentJob 注释
PropertiesJob 的属性
CurrentOffsetJob 当前处理完成的 Offset。只有 ExecuteType=Streaming 才有值
EndOffsetJob 获取到数据源端最大的 EndOffset。只有 ExecuteType=Streaming 才有值
LoadStatisticJob 的统计信息
ErrorMsgJob 执行的错误信息
JobRuntimeMsgJob 运行时的一些提示信息
LagBytes数据源端日志(MySQL Binlog / PostgreSQL WAL)的积压字节数,-1 表示当前不可用(例如 S3 数据源或全量快照阶段)。自 4.1.4 版本起,该列由原来的 Lag(单位:秒)改名为 LagBytes(单位:字节)
LastSourceEventTimestamp已提交 Offset 中记录的数据源端最新事件时间戳(Unix 秒),为空表示不可用。自 4.1.4 版本起新增
LastTaskSuccessTime最近一次 Task 成功完成的时间

查看 Task 状态

按 Job ID 查询其下所有子任务:

select * from tasks("type"="insert") where jobId='<job_id>';

返回结果列说明:

结果列说明
TaskId任务 ID
JobIDJobID
JobNameJob 名称
LabelTask 导入的 Label
StatusTask 的状态
ErrorMsgTask 失败信息
CreateTimeTask 的创建时间
StartTimeTask 的开始时间
FinishTimeTask 的完成时间
LoadStatisticTask 的统计信息
UserTask 的执行者
RunningOffset当前 Task 同步的 Offset 信息。只有 Job.ExecuteType=Streaming 才有值

暂停导入作业

手动暂停指定作业(暂停后不会被 autoResume 自动唤醒):

PAUSE JOB WHERE jobname = <job_name>;

恢复导入作业

恢复处于 PAUSED 状态的作业:

RESUME JOB WHERE jobName = <job_name>;

删除导入作业

彻底删除指定作业,删除后不再调度:

DROP JOB WHERE jobName = <job_name>;

通用参数

FE 配置参数

参数默认值说明
max_streaming_job_num1024最大的 Streaming 作业数量
job_streaming_task_exec_thread_num10用于执行 StreamingTask 的线程数
max_streaming_task_show_count100StreamingTask 在内存中最多保留的 task 执行记录

Job 通用导入配置参数

参数默认值说明
max_interval10当上游没有新增数据时,空闲的调度间隔,单位为秒。只接受整数(秒数),如 10;不支持带单位后缀(如 10s)。取值需 >= 1。

使用限制

同步范围、自动建表、一致性语义(exactly-once / at-least-once)已在能力对比中说明。本节只列出不支持的约束与行为。

主键表

仅支持同步带主键的上游表(两种同步方式均如此)。对应的 Doris 表为 Unique Key 表——自动建表同步下由 Doris 自动建成,SQL 映射同步下由你自己创建。无主键的表不支持。

Schema Change(DDL)

DDL 同步仅适用于自动建表同步;SQL 映射(TVF)不同步任何 DDL——cdc_stream() 表函数会强制把 schema_change_enabled 设为 false

  • PostgreSQL(4.1 起支持):仅 ADD COLUMNDROP COLUMN 会同步。列类型变更、RENAME COLUMN、约束 / 索引 / 分区变更不会同步——需在 Doris 端手工处理。
  • MySQL(自 4.1.4 起支持):仅 ADD COLUMNDROP COLUMN 会同步。列类型变更、RENAME COLUMN、约束 / 索引 / 分区变更不会同步——需在 Doris 端手工处理。

可以通过 Job 属性 schema_change_enabled 关闭该能力:

参数适用数据源默认值说明
schema_change_enabledMySQL、PostgreSQLtrue是否自动同步上游的 ADD COLUMN / DROP COLUMN。自 4.1.4 版本起支持
版本行为变更(4.1.4)
  • 自 4.1.4 版本起,同步新增列时不再传递上游列的 DEFAULT(MySQL 与 PostgreSQL 均如此)。新增列在 Doris 端不带默认值,历史数据不会回填。
  • PostgreSQL 的 Schema Change 检测改为基于 Relation 事件驱动:只识别 ADD COLUMNDROP COLUMN;同时包含新增和删除列的变更(可能是 RENAME)以及列类型变更会被跳过。该能力只在"自动建表同步(at-least-once)"路径上生效,TVF / exactly-once 路径不支持。

FAQ

MySQL 连接报错 Public Key Retrieval is not allowed

原因: 配置的 MySQL 用户使用 SHA256 密码认证方式,需要通过 TLS 等协议传输密码。

解决方案一: 在 JDBC URL 中添加 allowPublicKeyRetrieval=true 参数:

jdbc:mysql://127.0.0.1:3306?allowPublicKeyRetrieval=true

解决方案二: 将 MySQL 用户的认证方式改为 mysql_native_password

ALTER USER 'username'@'%' IDENTIFIED WITH mysql_native_password BY 'password';
FLUSH PRIVILEGES;