首批通过分布式安全可靠测评,为关键业务系统打造
使用 Flink DirectLoad 实现 OceanBase 旁路导入
更新时间:2026-04-09 14:12:04
概述
Apache Flink 是一个开源的分布式流处理框架,广泛应用于实时和批量数据处理场景。OceanBase Flink DirectLoad 连接器(flink-connector-oceanbase-directload)是专为 OceanBase 数据库设计的高性能数据导入工具,它利用 OceanBase 数据库的旁路导入技术,能够以极高的吞吐量将大批量数据快速写入数据库。
旁路导入 是 OceanBase 提供的一种高性能数据导入方式。与传统的 SQL 语句 INSERT 方式不同,旁路导入绕过 SQL 层的解析和事务处理流程,直接将数据写入存储层,从而实现更高的导入性能。这种方式特别适合需要一次性导入大量数据的场景。
DirectLoad 连接器与普通 JDBC 连接器对比
| 对比项 | DirectLoad 连接器 | JDBC 连接器 |
|---|---|---|
| 写入方式 | 旁路导入,直接写入存储层 | 通过 SQL 语句 INSERT 写入 |
| 吞吐量 | 极高,适合大批量数据 | 相对较低 |
| 适用场景 | 批量数据导入、数据迁移 | 实时流式写入、CDC 同步 |
| 数据流类型 | 仅支持有界流(Bounded) | 支持无界流(Unbounded) |
| 导入期间表状态 | 目标表被锁定,仅支持查询 | 表可正常读写 |
| 支持的操作 | INSERT、UPDATE_AFTER |
INSERT、UPDATE、DELETE |
核心优势
- 超高吞吐量:相比传统 JDBC 方式,性能提升数倍甚至数十倍。
- 并行写入:支持多节点并行写入,充分利用集群资源。
- 简单易用:通过 Flink SQL 即可使用,无需复杂配置。
- 灵活的冲突处理:支持多种主键冲突处理策略(停止、替换、忽略)。
使用说明与限制
在使用 DirectLoad 连接器前,请务必了解以下特性和限制。
支持的功能
- 批量数据写入:专为大批量数据导入设计,支持千万级、亿级数据快速写入。
- 有界流处理:支持有界数据源(如文件、数据库快照等)。
- 多种导入模式:支持
full(全量)、inc(增量)、inc_replace(增量替换)。 - 灵活的冲突策略:提供
STOP_ON_DUP、REPLACE、IGNORE三种主键冲突处理策略。 - 并行写入:支持 Flink 多并行度写入,充分利用集群资源。
- 兼容性:支持 OceanBase MySQL 模式和 Oracle 模式。
使用限制
- 仅支持有界流:数据源必须是有界的(Bounded Stream),不支持持续不断的无界流写入
- 导入期间锁表:旁路导入执行期间,目标表会被锁定,此时仅允许
SELECT查询操作,不允许对该表进行INSERT、UPDATE、DELETE等写入操作 - 不支持 DELETE 操作:Flink Table/SQL 侧仅支持
INSERT与UPDATE_AFTER变更类型写入,不支持DELETE、UPDATE_BEFORE - 批处理场景:推荐使用 Flink Batch 模式,不适合实时流式写入场景
注意
如果您需要实时流式写入或处理无界数据流,请使用项目中的 flink-connector-oceanbase 连接器,它基于 JDBC 实现,支持无界流和实时写入。
适用场景
DirectLoad 连接器适合以下场景:
- 数据迁移:从其他数据库或数据仓库批量迁移数据到 OceanBase 数据库。
- 离线 ETL:定期批量处理和导入数据,如每日数据汇总。
- 历史数据导入:一次性导入大量历史数据。
- 数据初始化:为新系统批量初始化数据。
不适用场景
以下场景不建议使用 DirectLoad 连接器:
- 实时流式写入:需要持续不断写入数据的实时场景。
- CDC 数据同步:需要实时同步数据库变更的场景。
- 低延迟要求:对写入延迟有严格要求的场景。
- 高频小批量写入:频繁的小批量数据写入。
版本要求与依赖
| 软件 | 版本要求 | 说明 |
|---|---|---|
| Apache Flink | 1.15 或更高版本 | 推荐使用 1.15+ 以获得更好的性能和稳定性 |
| JDK | 8 或更高版本 | — |
| OceanBase | 满足以下任一版本范围: • 4.3.0 BP1 及以上 • 4.2.4 - 4.3.0 • 4.2.1 BP7 - 4.2.2.0 |
支持 MySQL 模式和 Oracle 模式 |
| OceanBase (增量模式) | 4.3.2 或更高版本 | 使用 inc 或 inc_replace 导入模式时需要 |
快速开始
本节将指导您快速上手使用 DirectLoad 连接器,从获取连接器到完成第一次数据导入。
步骤一:快速部署本地 Flink 集群(单机版)
下载 Apache Flink
前往 Flink 官方下载页面,选择 Stable Release(推荐 1.15+,如 1.18 或 1.19)。
例如(Linux / macOS 终端命令):
# 下载 Flink(以 1.18.1 为例)
wget https://archive.apache.org/dist/flink/flink-1.19.3/flink-1.19.3-bin-scala_2.12.tgz
# 解压
tar -xzf flink-1.19.3-bin-scala_2.12.tgz
# 进入目录
cd flink-1.19.3
# 设置环境变量(可选但推荐)
export FLINK_HOME=$(pwd)
启动本地 Flink 集群
在 Flink 目录下执行:
# 启动集群(包含 JobManager + TaskManager)
$FLINK_HOME/bin/start-cluster.sh
成功启动后,你会看到类似输出:
Starting cluster.
Starting standalonesession daemon on host your-hostname.
Starting taskexecutor daemon on host your-hostname.
如果启动失败,可以检查日志:
# 查看 JobManager 日志
tail -f $FLINK_HOME/log/flink-*-standalonesession-*.log
# 查看 TaskManager 日志
tail -f $FLINK_HOME/log/flink-*-taskexecutor-*.log
验证 Flink 是否正常运行
打开浏览器访问:http://localhost:8081, 您可以看到 Flink Dashboard,显示 1 个 TaskManager,可用 slots ≥ 1。
步骤二:部署 OceanBase DirectLoad 连接器 JAR
下载连接器 JAR
前往 Maven Central:https://repo1.maven.org/maven2/com/oceanbase/flink-sql-connector-oceanbase-directload/
选择版本(例如 1.5.0),下载对应的 JAR 文件:
# 示例:下载 1.5.0 版本
wget https://repo1.maven.org/maven2/com/oceanbase/flink-sql-connector-oceanbase-directload/1.0.0/flink-sql-connector-oceanbase-directload-1.5.0.jar
注意:文件名必须是 flink-sql-connector-oceanbase-directload-<version>.jar,这是 Flink SQL 识别连接器的关键。
将 JAR 放入 Flink 的 lib/ 目录
# 复制 JAR 到 Flink lib 目录
cp flink-sql-connector-oceanbase-directload-1.5.0.jar $FLINK_HOME/lib/
重要提示:Flink 在启动时会自动加载 lib/ 目录下的所有 JAR 包,因此你的连接器会被自动注册。
重启 Flink 集群(使 JAR 生效)
# 先停止
$FLINK_HOME/bin/stop-cluster.sh
# 再启动
$FLINK_HOME/bin/start-cluster.sh
步骤三:获取数据库连接信息
联系 OceanBase 数据库部署人员或者管理员获取相应的数据库连接串,例如:
obclient -h$host -P$port -u$user_name -p$password -D$database_name
参数说明:
$host:提供 OceanBase 数据库连接 IP。OceanBase 数据库代理(OceanBase Database Proxy,ODP)连接方式使用的是一个 ODP 地址;直连方式使用的是 OBServer 节点的 IP 地址。$port:提供 OceanBase 数据库连接端口。ODP 连接的方式默认是2883,在部署 ODP 时可自定义;直连方式默认是2881,在部署 OceanBase 数据库时可自定义。$database_name:需要访问的数据库名称。注意
连接租户的用户需要拥有该数据库的
CREATE、INSERT、DROP和SELECT权限。更多有关用户权限的信息,请参见 MySQL 模式下的权限分类。$user_name:提供租户的连接账户。ODP 连接的常用格式:用户名@租户名#集群名或者集群名:租户名:用户名;直连方式格式:用户名@租户名。$password:提供账户密码。
更多连接串的信息,请参见 通过 OBClient 连接 OceanBase 租户。
注意
连接 OBServer 服务端时,sys 租户下查询系统视图 DBA_OB_SERVERS 即可获取 OBServer 的 RPC 端口号。
确认网络连通性
确保运行 Flink 作业的环境能够访问 OceanBase 数据库:
测试命令:
# 测试 RPC 端口连通性
nc -zv <oceanbase-host> <rpc-port>
步骤四:在 OceanBase 中创建目标表
示例:
-- 连接到 OceanBase 数据库
USE test;
-- 创建目标表
CREATE TABLE `t_user` (
`id` INT NOT NULL,
`username` VARCHAR(50) DEFAULT NULL,
`age` INT DEFAULT NULL,
`score` DECIMAL(10,2) DEFAULT NULL,
PRIMARY KEY (`id`)
) COMMENT '用户信息表';
步骤五:启动 Flink SQL Client 并测试
# 启动 SQL Client
$FLINK_HOME/bin/sql-client.sh
成功启动后,您将看到 Flink SQL Client 的交互式命令行界面:
▒▓██▓██▒
▓████▒▒█▓▒▓███▓▒
▓███▓░░ ▒▒▒▓██▒ ▒
░██▒ ▒▒▓▓█▓▓▒░ ▒████
██▒ ░▒▓███▒ ▒█▒█▒
░▓█ ███ ▓░▒██
▓█ ▒▒▒▒▒▓██▓░▒░▓▓█
█░ █ ▒▒░ ███▓▓█ ▒█▒▒▒
████░ ▒▓█▓ ██▒▒▒ ▓███▒
░▒█▓▓██ ▓█▒ ▓█▒▓██▓ ░█░
▓░▒▓████▒ ██ ▒█ █▓░▒█▒░▒█▒
███▓░██▓ ▓█ █ █▓ ▒▓█▓▓█▒
░██▓ ░█░ █ █▒ ▒█████▓▒ ██▓░▒
███░ ░ █░ ▓ ░█ █████▒░░ ░█░▓ ▓░
██▓█ ▒▒▓▒ ▓███████▓░ ▒█▒ ▒▓ ▓██▓
▒██▓ ▓█ █▓█ ░▒█████▓▓▒░ ██▒▒ █ ▒ ▓█▒
▓█▓ ▓█ ██▓ ░▓▓▓▓▓▓▓▒ ▒██▓ ░█▒
▓█ █ ▓███▓▒░ ░▓▓▓███▓ ░▒░ ▓█
██▓ ██▒ ░▒▓▓███▓▓▓▓▓██████▓▒ ▓███ █
▓███▒ ███ ░▓▓▒░░ ░▓████▓░ ░▒▓▒ █▓
█▓▒▒▓▓██ ░▒▒░░░▒▒▒▒▓██▓░ █▓
██ ▓░▒█ ▓▓▓▓▒░░ ▒█▓ ▒▓▓██▓ ▓▒ ▒▒▓
▓█▓ ▓▒█ █▓░ ░▒▓▓██▒ ░▓█▒ ▒▒▒░▒▒▓█████▒
██░ ▓█▒█▒ ▒▓▓▒ ▓█ █░ ░░░░ ░█▒
▓█ ▒█▓ ░ █░ ▒█ █▓
█▓ ██ █░ ▓▓ ▒█▓▓▓▒█░
█▓ ░▓██░ ▓▒ ▓█▓▒░░░▒▓█░ ▒█
██ ▓█▓░ ▒ ░▒█▒██▒ ▓▓
▓█▒ ▒█▓▒░ ▒▒ █▒█▓▒▒░░▒██
░██▒ ▒▓▓▒ ▓██▓▒█▒ ░▓▓▓▓▒█▓
░▓██▒ ▓░ ▒█▓█ ░░▒▒▒
▒▓▓▓▓▓▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒░░▓▓ ▓░▒█░
______ _ _ _ _____ ____ _ _____ _ _ _ BETA
| ____| (_) | | / ____|/ __ \| | / ____| (_) | |
| |__ | |_ _ __ | | __ | (___ | | | | | | | | |_ ___ _ __ | |_
| __| | | | '_ \| |/ / \___ \| | | | | | | | | |/ _ \ '_ \| __|
| | | | | | | | < ____) | |__| | |____ | |____| | | __/ | | | |_
|_| |_|_|_| |_|_|\_\ |_____/ \___\_\______| \_____|_|_|\___|_| |_|\__|
Welcome! Enter 'HELP;' to list all available commands. 'QUIT;' to exit.
Flink SQL>
步骤六:第一个 Flink SQL 示例
启动 Flink SQL Client
cd $FLINK_HOME
./bin/sql-client.sh
设置为 Batch 模式
-- 可选:建议设置运行模式为 BATCH,可获得更好的性能
SET 'execution.runtime-mode' = 'BATCH';
-- 可选:设置并行度(根据数据量和集群资源调整并行度)
SET 'parallelism.default' = '4';
创建 DirectLoad Sink 表
在 Flink SQL 中创建对应的目标表,映射到 OceanBase 的 t_user 表:
CREATE TABLE t_user (
id INT,
username STRING,
age INT,
score DECIMAL(10,2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = 'xxx.x.x.x', -- OceanBase 主机地址
'port' = 'xxxx', -- RPC 端口
'tenant-name' = 'your_tenant-name', -- 租户名
'username' = 'your_username', -- 用户名
'password' = 'your_password', -- 密码
'schema-name' = 'test', -- 数据库名
'table-name' = 't_user' -- 表名
);
参数说明:
connector:固定值 oceanbase-directloadhost和port:OceanBase 的主机地址和 RPC 端口号。连接 OBServer 服务端时,sys 租户下查询系统视图DBA_OB_SERVERS即可获取 OBServer 的 RPC 端口号。tenant-name:租户名称username和password:数据库用户凭证(注意:用户名不包含 @tenant 后缀)schema-name和table-name:目标数据库和表名
说明
以上使用了最小参数集。更多参数配置请参考后续章节中的 -配置参数详解。
插入测试数据
-- 插入几条测试数据
INSERT INTO t_user
VALUES
(1, 'Alice', 25, 95.5),
(2, 'Bob', 30, 88.0),
(3, 'Charlie', 28, 92.3),
(4, 'David', 35, 87.8),
(5, 'Eve', 22, 96.0);
执行后,Flink 将提交作业并开始导入数据,您会在控制台看到类似如下的输出:
[INFO] Submitting SQL update statement to the cluster...
[INFO] SQL update statement has been successfully submitted to the cluster:
Job ID: xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx
作业执行说明:
- 作业将在后台执行,您可以在 Flink Web UI(http://localhost:8081)中查看作业进度。
- 作业状态会经历:
INITIALIZING→RUNNING→FINISHED。 - DirectLoad 的最终提交发生在作业结束阶段,请务必等待作业状态变为
FINISHED后再验证数据。
验证数据写入
作业成功完成后,在 OceanBase 中查询验证:
SELECT * FROM test.t_user ORDER BY id;
查询结果如下:
+----+----------+------+-------+
| id | username | age | score |
+----+----------+------+-------+
| 1 | Alice | 25 | 95.50 |
| 2 | Bob | 30 | 88.00 |
| 3 | Charlie | 28 | 92.30 |
| 4 | David | 35 | 87.80 |
| 5 | Eve | 22 | 96.00 |
+----+----------+------+-------+
5 rows in set
步骤七:关闭和清理(可选)
完成测试后,您可以停止 Flink 集群:
# 停止 Flink 集群
cd $FLINK_HOME
./bin/stop-cluster.sh
在 SQL Client 中退出:
-- 在 SQL Client 中执行
QUIT;
或者直接按 Ctrl+D 退出 SQL Client。
重要提示
- 等待作业完成:旁路导入的最终提交发生在作业结束阶段,请务必等待作业状态变为
"FINISHED"后再验证数据。 - 导入期间锁表:在
INSERT语句执行期间,目标表t_user会被锁定,此时只能查询该表,无法执行写入操作。 - 推荐使用 Batch 模式:DirectLoad 连接器仅支持有界流,推荐使用 Batch 模式并确保数据源是有界的。
- 错误处理:如果作业失败,请检查日志中的错误信息,常见问题包括网络连接失败、账号权限不足、参数配置错误等。
配置参数详解
本节详细介绍 DirectLoad 连接器的所有配置参数。带 * 标记的参数为必填项。
连接参数
| 参数名 | 必填 | 默认值 | 类型 | 说明 |
|---|---|---|---|---|
connector* |
是 | 无 | String | 连接器类型,固定值:oceanbase-directload |
host* |
是 | 无 | String | OceanBase 数据库的主机地址,如 127.0.0.1 或域名 |
port* |
是 | 2882 |
Integer | 旁路导入使用的 RPC 端口,默认为 2882 |
tenant-name* |
是 | 无 | String | 租户名称,如 test |
username* |
是 | 无 | String | 数据库用户名,如 root。注意:这里填写纯用户名,而不是连接串格式(如 root@test) |
password* |
是 | 无 | String | 数据库密码 |
schema-name* |
是 | 无 | String | 数据库名(Schema),如 test |
table-name* |
是 | 无 | String | 目标表名,如 t_user |
导入行为参数
| 参数名 | 必填 | 默认值 | 类型 | 说明 |
|---|---|---|---|---|
load-method |
否 | full |
String | 导入模式,支持:full、inc、inc_replace详见下文 “ load-method 详解” |
dup-action |
否 | REPLACE |
String | 主键冲突处理策略,支持: - STOP_ON_DUP:遇到冲突立即停止导入- REPLACE:替换已存在的记录- IGNORE:忽略冲突记录详见下文 “ dup-action 详解” |
max-error-rows |
否 | 0 |
Long | 最大可容忍的错误行数。超过此阈值导入失败。 设置为 0 表示不容忍任何错误 |
性能调优参数
| 参数名 | 必填 | 默认值 | 类型 | 说明 |
|---|---|---|---|---|
parallel |
否 | 8 |
Integer | 服务端并行度。该参数决定了 OceanBase 服务端使用多少 CPU 资源来处理本次导入任务。 详见下文 “ parallel 详解” |
buffer-size |
否 | 1024 |
Integer | 写入缓冲区大小(单位:行数)。累计到该行数后触发一次 flush 写入。 详见下文 “ buffer-size 调优建议” |
timeout |
否 | 7d |
Duration | 单次旁路导入任务的超时时间。支持格式:1d、12h、30m、3600s |
heartbeat-timeout |
否 | 60s |
Duration | 客户端心跳超时时间 |
heartbeat-interval |
否 | 10s |
Duration | 客户端心跳间隔时间 |
核心参数详细说明
load-method 详解
load-method 参数决定了旁路导入的模式,不同模式适用于不同场景。
full(全量导入)
- 默认值,适合空表或仅含少量数据的表
- 导入的数据直接写入 major sstable 中
- 列存表优势:导入完成后数据即为列存格式,查询性能最优,无需额外合并
- 使用场景:
- 首次导入数据到空表
- 目标表数据量很小,可接受重建
- 追求最优的列存查询性能
-- full 模式示例
CREATE TABLE t_sink (...) WITH (
'connector' = 'oceanbase-directload',
'load-method' = 'full',
...
);
inc(普通增量导入)
- 适合向非空表追加数据
- 会进行主键冲突检查,如果发现冲突按
dup-action策略处理 - 版本要求:OceanBase 4.3.2 及以上版本
- 限制:暂不支持
dup-action为REPLACE - 列存表注意:数据写入转储(转储暂不支持列存),导入后查询性能为行存性能,需要等待一次合并才能达到列存查询性能
-- inc 模式示例
CREATE TABLE t_sink (...) WITH (
'connector' = 'oceanbase-directload',
'load-method' = 'inc',
'dup-action' = 'STOP_ON_DUP', -- 不能使用 REPLACE
...
);
inc_replace(增量替换导入)
- 适合向非空表导入数据,且需要覆盖已存在的记录
- 不进行主键冲突检查,直接覆盖原有主键的数据(相当于
REPLACE的效果) - 版本要求:OceanBase 4.3.2 及以上版本
dup-action参数会被忽略(因为不检查冲突)- 列存表注意:与
inc模式相同,需要合并后才能达到列存查询性能
-- inc_replace 模式示例
CREATE TABLE t_sink (...) WITH (
'connector' = 'oceanbase-directload',
'load-method' = 'inc_replace', -- dup-action 参数无效
...
);
模式选择建议
| 场景 | 推荐模式 | 理由 |
|---|---|---|
| 首次导入到空表 | full |
性能最优,列存表查询效果最佳 |
| 向非空表导入,可能有冲突且需要覆盖 | inc_replace |
自动覆盖旧数据 |
| 向非空表导入,不允许冲突 | inc + dup-action=STOP_ON_DUP |
严格控制数据一致性 |
dup-action 详解
dup-action 参数用于处理主键冲突的情况(仅在 load-method=full 或 inc 时生效)。
STOP_ON_DUP(停止导入)
- 遇到主键冲突时立即停止导入,整个作业失败
- 适用场景:
- 对数据一致性要求严格
- 不允许出现主键冲突
- 需要快速发现数据质量问题
'dup-action' = 'STOP_ON_DUP'
REPLACE(替换)
- 遇到主键冲突时,用新记录替换旧记录
- 适用场景:
- 允许覆盖已存在的数据
- 导入的是最新版本数据
- 注意:在
load-method=inc时暂不支持
'dup-action' = 'REPLACE'
IGNORE(忽略)
- 遇到主键冲突时,保留旧记录,忽略新记录
- 适用场景:
- 以旧数据为准
- 增量导入时避免覆盖已存在的记录
'dup-action' = 'IGNORE'
parallel 详解
parallel 参数是服务端并行度,它决定了 OceanBase 服务端使用多少 CPU 资源来处理本次导入任务。
重要特性
- 与客户端并发无关:
parallel是服务端参数,与 Flink 的并行度是两个独立的概念 - 受租户配置限制:服务端会根据租户 CPU 配置自动限制并行度上限,客户端设置超出范围不会报错
- 受分区分布影响:实际并行度还受表的分区分布影响
实际并行度计算规则
- 单节点并行度上限 =
MIN(租户CPU数 × 2, parallel 配置值) - 实际总并行度 =
单节点并行度上限 × 分区分布的节点数
示例 1:单节点场景
- 租户配置:2C
parallel设置:10- 表的分区在 1 个节点
- 实际并行度 =
MIN(2 × 2, 10) × 1 = 4
示例 2:多节点场景
- 租户配置:2C
parallel设置:10- 表的分区均匀分布在 2 个节点
- 实际并行度 =
MIN(2 × 2, 10) × 2 = 8
调优建议
- 对于大数据量导入,适当调大
parallel可显著缩短 commit 阶段耗时 - 建议根据租户 CPU 配置设置合理值,过大无效,过小影响性能
- 一般设置为租户 CPU 数的 2–4 倍是合理的起点
-- 示例:4C 租户建议设置
'parallel' = '8' -- 或 '16'
buffer-size 调优建议
buffer-size 表示客户端写入缓冲区大小,单位是行数。累计到该行数后触发一次 flush 写入。
调优建议
- 数据量大且单行小:可适当调大(如 2048、4096),减少 flush 次数,提升吞吐
- 单行数据大:避免设置过大,防止内存压力,建议保持默认值或适当调小
- 内存紧张:减小
buffer-size,避免 TaskManager OOM
推荐值
| 场景 | 推荐值 |
|---|---|
| 单行 < 1KB,内存充足 | 2048 – 4096 |
| 单行 1KB – 10KB | 1024(默认值) |
| 单行 > 10KB 或内存紧张 | 512 |
-- 示例
'buffer-size' = '2048'
使用示例
本节提供更完整和实用的使用示例,涵盖不同场景和配置。
示例 1:基础批量导入
这是最基础的使用示例,适合初学者快速上手。
-- 1. 设置为 Batch 模式
SET 'execution.runtime-mode' = 'BATCH';
-- 2. 创建 Sink 表(使用最小参数集)
CREATE TABLE orders_sink (
order_id BIGINT,
user_id BIGINT,
product_name STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = '127.0.0.1',
'port' = '2882',
'tenant-name' = 'test',
'username' = 'root',
'password' = 'your_password',
'schema-name' = 'test',
'table-name' = 'orders'
);
-- 3. 从数据源导入(这里假设从另一个表读取)
INSERT INTO orders_sink
SELECT order_id, user_id, product_name, amount, order_time
FROM orders_source;
示例 2:全量导入到空表(full 模式)
适用于首次导入大量数据到空表的场景。
SET 'execution.runtime-mode' = 'BATCH';
SET 'parallelism.default' = '8'; -- 根据数据量调整
CREATE TABLE user_profile_sink (
user_id BIGINT,
username STRING,
email STRING,
age INT,
city STRING,
register_time TIMESTAMP(3),
PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = '127.0.0.1',
'port' = '2882',
'tenant-name' = 'prod',
'username' = 'admin',
'password' = 'your_password',
'schema-name' = 'userdb',
'table-name' = 'user_profile',
-- 使用 full 模式(适合空表)
'load-method' = 'full',
-- 性能调优
'parallel' = '16', -- 服务端并行度
'buffer-size' = '2048'
);
-- 从 CSV 文件或其他数据源导入
INSERT INTO user_profile_sink
SELECT * FROM user_profile_csv_source;
示例 3:增量替换导入(inc_replace 模式)
适用于需要用新数据覆盖旧数据的场景。
SET 'execution.runtime-mode' = 'BATCH';
CREATE TABLE product_info_sink (
product_id BIGINT,
product_name STRING,
price DECIMAL(10, 2),
stock INT,
update_time TIMESTAMP(3),
PRIMARY KEY (product_id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = '127.0.0.1',
'port' = '2882',
'tenant-name' = 'ecommerce',
'username' = 'root',
'password' = 'your_password',
'schema-name' = 'productdb',
'table-name' = 'product_info',
-- 使用 inc_replace 模式(自动覆盖已存在的数据)
'load-method' = 'inc_replace', -- dup-action 参数会被忽略
-- 性能配置
'parallel' = '16',
'buffer-size' = '2048'
);
-- 导入更新数据(自动覆盖旧记录)
INSERT INTO product_info_sink
SELECT * FROM product_info_latest;
示例 4:从 Kafka 读取有界数据并导入
这个示例展示了如何从 Kafka 读取有界数据(通过 bounded 模式)并导入到 OceanBase。
SET 'execution.runtime-mode' = 'BATCH';
-- 创建 Kafka Source(有界模式)
CREATE TABLE kafka_source (
event_id BIGINT,
event_type STRING,
user_id BIGINT,
event_data STRING,
event_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'user-events',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'flink-import-group',
'scan.startup.mode' = 'earliest-offset',
'scan.bounded.mode' = 'latest-offset', -- 有界模式:读取到最新offset后停止
'format' = 'json'
);
-- 创建 OceanBase Sink
CREATE TABLE event_sink (
event_id BIGINT,
event_type STRING,
user_id BIGINT,
event_data STRING,
event_time TIMESTAMP(3),
PRIMARY KEY (event_id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = '127.0.0.1',
'port' = '2882',
'tenant-name' = 'analytics',
'username' = 'root',
'password' = 'your_password',
'schema-name' = 'eventdb',
'table-name' = 'events',
'load-method' = 'full',
'parallel' = '16'
);
-- 导入数据
INSERT INTO event_sink
SELECT * FROM kafka_source;
示例 5:从文件导入数据
从 CSV 文件批量导入数据到 OceanBase。
SET 'execution.runtime-mode' = 'BATCH';
-- 创建文件 Source
CREATE TABLE csv_source (
id BIGINT,
name STRING,
age INT,
salary DECIMAL(10, 2)
) WITH (
'connector' = 'filesystem',
'path' = 'file:///data/employees.csv',
'format' = 'csv',
'csv.field-delimiter' = ',',
'csv.ignore-parse-errors' = 'false'
);
-- 创建 OceanBase Sink
CREATE TABLE employee_sink (
id BIGINT,
name STRING,
age INT,
salary DECIMAL(10, 2),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase-directload',
'host' = '192.168.1.100',
'port' = '2882',
'tenant-name' = 'hrms',
'username' = 'admin',
'password' = 'your_password',
'schema-name' = 'hr',
'table-name' = 'employees',
'load-method' = 'full',
'parallel' = '8',
'buffer-size' = '2048',
'max-error-rows' = '10' -- 允许最多10行错误
);
-- 执行导入
INSERT INTO employee_sink
SELECT * FROM csv_source;
使用 DirectLoad 连接器的最佳实践
性能优化
合理设置 Flink 并行度
Flink 并行度决定了有多少个并行的 Writer 任务。合理设置可以充分利用集群资源。
-- 根据数据量和集群资源设置并行度
SET 'parallelism.default' = '8';
调优 parallel 参数(服务端并行度)
parallel 参数对 commit 阶段性能影响很大。
调优策略:
- 大数据量导入时,适当调大
parallel可大幅缩短 commit 阶段耗时 - 建议设置为租户 CPU 数的 2–4 倍
- 不用担心设置过大,服务端会自动限制
-- 示例:8C 租户
'parallel' = '16' -- 或 '32'
性能对比示例
parallel |
commit 阶段耗时 |
|---|---|
| 8(默认) | 约 10 分钟 |
| 16 | 约 5 分钟 |
| 32 | 约 3 分钟 |
调优 buffer-size
根据数据特征调整缓冲区大小:
-- 小行数据(每行 < 1KB)
'buffer-size' = '4096'
-- 中等行数据(每行 1–10KB)
'buffer-size' = '1024' -- 默认值
-- 大行数据(每行 > 10KB)
'buffer-size' = '512'
生产环境建议
选择合适的导入窗口期
由于旁路导入期间会锁表,建议:
- 选择业务低峰期:如凌晨、周末
- 提前通知相关方:避免影响其他业务
- 设置合理的超时时间:防止导入任务长时间占用表
-- 设置2小时超时(根据数据量评估)
'timeout' = '2h'
根据场景选择 load-method
| 场景 | 推荐 load-method |
原因 |
|---|---|---|
| 首次导入到空表 | full |
性能最优,列存查询效果最佳 |
| 定期追加增量数据 | inc |
适合增量场景,支持冲突检查 |
| 每日全量更新维度表 | inc_replace |
自动覆盖旧数据,无需手动删除 |
| 历史数据回溯 | full(导入到临时表后切换) |
避免影响在线表 |
列存表使用建议
如果目标是列存表,需要注意不同 load-method 的影响:
full模式(推荐):- 数据直接写为列存格式
- 导入完成后查询性能最优
- 仅适合空表或可接受重建的表
inc/inc_replace模式:- 数据写入转储(行存格式)
- 导入后查询性能为行存性能
- 需要等待一次 major compaction 后才能达到列存性能
- 如果追求高查询性能,建议手动触发合并:
-- 在 OceanBase 中手动触发合并
ALTER SYSTEM MAJOR FREEZE;
常见问题
Q1:导入期间其他写入操作失败怎么办?
问题描述:在执行旁路导入时,尝试对目标表进行 INSERT/UPDATE/DELETE 操作失败。
原因:这是旁路导入的固有特性。导入期间目标表会被锁定,仅允许 SELECT 操作。
解决方案:
方案一:错峰导入
- 将导入任务安排在业务低峰期(如凌晨)
- 提前与业务方沟通,暂停写入操作
方案二:使用中间表
-- 1. 导入到临时表 CREATE TABLE target_table_tmp LIKE target_table; -- 使用 DirectLoad 导入到 target_table_tmp -- 2. 切换表名 RENAME TABLE target_table TO target_table_old, target_table_tmp TO target_table;- 先导入到临时表
- 导入完成后通过表切换或数据合并操作
方案三:使用普通 JDBC 连接器
- 如果无法接受锁表,改用
flink-connector-oceanbase连接器。 - 牺牲部分性能换取表的可用性。
- 如果无法接受锁表,改用
Q2:作业一直不结束 / 数据未提交怎么办?
问题描述:Flink 作业一直处于 RUNNING 状态,OceanBase 中查询不到数据。
原因:DirectLoad 连接器的最终 commit 发生在输入结束(end-of-input)阶段。如果输入是无界流,作业不会结束,数据也不会提交。
排查步骤:
检查是否设置了 Batch 模式:
SET 'execution.runtime-mode' = 'BATCH';检查数据源是否为有界流:
- 文件数据源:天然有界
- Kafka:需设置
scan.bounded.mode - JDBC:天然有界
- 自定义 Source:检查是否正确发送 end-of-input 信号
查看 Flink Web UI:
- 检查作业状态
- 查看是否有 Backpressure
- 检查各个算子的处理进度
Q3:如何选择 load-method?
决策树:
是否为首次导入到空表?
├─ 是 → 使用 full 模式(性能最优,列存效果最佳)
└─ 否(向非空表导入)
└─ 是否需要覆盖已存在的记录?
├─ 是 → 使用 inc_replace 模式
└─ 否 → 使用 inc 模式
└─ 是否允许主键冲突?
├─ 否 → dup-action = STOP_ON_DUP
├─ 允许(保留新数据) → dup-action = REPLACE
└─ 允许(保留旧数据) → dup-action = IGNORE
特殊场景:
- 列存表追求最优查询性能:优先使用
full模式 - OceanBase 版本 < 4.3.2:只能使用
full模式 - 需要频繁增量导入:使用
inc或inc_replace模式
Q4:列存表使用 inc/inc_replace 模式后查询很慢怎么办?
问题描述:使用 inc 或 inc_replace 模式向列存表导入数据后,查询性能不如预期。
原因:inc/inc_replace 模式的数据写入转储(行存格式),需要等待一次 major compaction 后才能转为列存格式。
解决方案:
手动触发 major compaction(推荐):
-- 在 OceanBase 中执行 ALTER SYSTEM MAJOR FREEZE; -- 检查合并进度 SELECT * FROM oceanbase.DBA_OB_MAJOR_COMPACTION;等待自动合并:
OceanBase 会定期自动触发合并,但可能需要等待较长时间(取决于配置)。
Q5:导入速度很慢怎么优化?
可能原因及解决方案:
Flink 并行度过低:
SET 'parallelism.default' = '16'; -- 增加并行度服务端并行度过低:
'parallel' = '16' -- 增加服务端并行度buffer-size 过小:
'buffer-size' = '2048' -- 增大缓冲区OceanBase 资源不足:
- 检查 OceanBase 集群负载
- 考虑扩容或选择低峰期导入