---
title: "使用 Flink DirectLoad 实现 OceanBase 旁路导入 - OceanBase 数据库 V4.3.5 | OceanBase 文档中心"
description: 使用 Flink DirectLoad 实现 OceanBase 旁路导入 概述 Apache Flink 是一个开源的分布式流处理框架，广泛应用于实时和批量数据处理场景。 OceanBase Flink DirectLoad 连接器 （ flink-connector-oceanbase-directload ）是…
---
切换语言

- 中文站 - 简体中文
- International - English
- 日本站 - 日本語

文档反馈![](https://mdn.alipayobjects.com/huamei_22khvb/afts/img/A*P8CuR4UJ_FkAAAAAAAAAAAAADiGDAQ/original) OceanBase 数据库分布式版 - V 4.3.5 LTS

# 使用 Flink DirectLoad 实现 OceanBase 旁路导入

更新时间：2026-04-09 14:12:04

[编辑](https://github.com/oceanbase/oceanbase-doc/edit/V4.3.5/zh-CN/680.ecological-integration/400.data-ingestion/100.flink/300.load-data-bypass-through-flink.md)  

## 概述

Apache Flink 是一个开源的分布式流处理框架，广泛应用于实时和批量数据处理场景。**OceanBase Flink DirectLoad 连接器**（[`flink-connector-oceanbase-directload`](https://github.com/oceanbase/flink-connector-oceanbase/releases)）是专为 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 官方下载页面](https://flink.apache.org/downloads.html)，选择 **Stable Release**（推荐 1.15+，如 1.18 或 1.19）。

例如（Linux / macOS 终端命令）：

```bash
# 下载 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 目录下执行：

```bash
# 启动集群（包含 JobManager + TaskManager）
$FLINK_HOME/bin/start-cluster.sh

```

成功启动后，你会看到类似输出：

```shell
Starting cluster.
Starting standalonesession daemon on host your-hostname.
Starting taskexecutor daemon on host your-hostname.

```

如果启动失败，可以检查日志：

```shell
# 查看 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 文件：

```bash
# 示例：下载 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/` 目录

```bash
# 复制 JAR 到 Flink lib 目录
cp flink-sql-connector-oceanbase-directload-1.5.0.jar $FLINK_HOME/lib/

```

**重要提示**：Flink 在启动时会自动加载 `lib/` 目录下的所有 JAR 包，因此你的连接器会被自动注册。

#### 重启 Flink 集群（使 JAR 生效）

```bash
# 先停止
$FLINK_HOME/bin/stop-cluster.sh

# 再启动
$FLINK_HOME/bin/start-cluster.sh

```

### 步骤三：获取数据库连接信息

联系 OceanBase 数据库部署人员或者管理员获取相应的数据库连接串，例如：

```shell
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 模式下的权限分类](https://www.oceanbase.com/docs/common-oceanbase-database-cn-1000000002016100)。
 - `$user_name`：提供租户的连接账户。ODP 连接的常用格式：`用户名@租户名#集群名` 或者 `集群名:租户名:用户名`；直连方式格式：`用户名@租户名`。
 - `$password`：提供账户密码。

更多连接串的信息，请参见 [通过 OBClient 连接 OceanBase 租户](https://www.oceanbase.com/docs/common-oceanbase-database-cn-1000000002013251)。

#### 注意

连接 OBServer 服务端时，sys 租户下查询系统视图 `DBA_OB_SERVERS` 即可获取 OBServer 的 RPC 端口号。

#### 确认网络连通性

确保运行 Flink 作业的环境能够访问 OceanBase 数据库：

测试命令：

```bash
# 测试 RPC 端口连通性
nc -zv <oceanbase-host> <rpc-port>

```

### 步骤四：在 OceanBase 中创建目标表

示例：

```sql
-- 连接到 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 并测试

```bash
# 启动 SQL Client
$FLINK_HOME/bin/sql-client.sh

```

成功启动后，您将看到 Flink SQL Client 的交互式命令行界面：

```sql
▒▓██▓██▒
▓████▒▒█▓▒▓███▓▒
▓███▓░░        ▒▒▒▓██▒  ▒
░██▒   ▒▒▓▓█▓▓▒░      ▒████
██▒         ░▒▓███▒    ▒█▒█▒
░▓█            ███   ▓░▒██
▓█       ▒▒▒▒▒▓██▓░▒░▓▓█
█░ █   ▒▒░       ███▓▓█ ▒█▒▒▒
████░   ▒▓█▓      ██▒▒▒ ▓███▒
░▒█▓▓██       ▓█▒    ▓█▒▓██▓ ░█░
▓░▒▓████▒ ██         ▒█    █▓░▒█▒░▒█▒
███▓░██▓  ▓█           █   █▓ ▒▓█▓▓█▒
░██▓  ░█░            █  █▒ ▒█████▓▒ ██▓░▒
███░ ░ █░          ▓ ░█ █████▒░░    ░█░▓  ▓░
██▓█ ▒▒▓▒          ▓███████▓░       ▒█▒ ▒▓ ▓██▓
▒██▓ ▓█ █▓█       ░▒█████▓▓▒░         ██▒▒  █ ▒  ▓█▒
▓█▓  ▓█ ██▓ ░▓▓▓▓▓▓▓▒              ▒██▓           ░█▒
▓█    █ ▓███▓▒░              ░▓▓▓███▓          ░▒░ ▓█
██▓    ██▒    ░▒▓▓███▓▓▓▓▓██████▓▒            ▓███  █
▓███▒ ███   ░▓▓▒░░   ░▓████▓░                  ░▒▓▒  █▓
█▓▒▒▓▓██  ░▒▒░░░▒▒▒▒▓██▓░                            █▓
██ ▓░▒█   ▓▓▓▓▒░░  ▒█▓       ▒▓▓██▓    ▓▒          ▒▒▓
▓█▓ ▓▒█  █▓░  ░▒▓▓██▒            ░▓█▒   ▒▒▒░▒▒▓█████▒
██░ ▓█▒█▒  ▒▓▓▒  ▓█                █░      ░░░░   ░█▒
▓█   ▒█▓   ░     █░                ▒█              █▓
█▓   ██         █░                 ▓▓        ▒█▓▓▓▒█░
█▓ ░▓██░       ▓▒                  ▓█▓▒░░░▒▓█░    ▒█
██   ▓█▓░      ▒                    ░▒█▒██▒      ▓▓
▓█▒   ▒█▓▒░                         ▒▒ █▒█▓▒▒░░▒██
░██▒    ▒▓▓▒                     ▓██▓▒█▒ ░▓▓▓▓▒█▓
░▓██▒                          ▓░  ▒█▓█  ░░▒▒▒
▒▓▓▓▓▓▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒▒░░▓▓  ▓░▒█░

______ _ _       _       _____  ____  _         _____ _ _            _  BETA
|  ____| (_)     | |     / ____|/ __ \| |       / ____| (_)          | |
| |__  | |_ _ __ | | __ | (___ | |  | | |      | |    | |_  ___ _ __ | |_
|  __| | | | '_ \| |/ /  \___ \| |  | | |      | |    | | |/ _ \ '_ \| __|
| |    | | | | | |   <   ____) | |__| | |____  | |____| | |  __/ | | | |_
|_|    |_|_|_| |_|_|\_\ |_____/ \___\_\______|  \_____|_|_|\___|_| |_|\__|

Welcome! Enter 'HELP;' to list all available commands. 'QUIT;' to exit.

Flink SQL>

```

### 步骤六：第一个 Flink SQL 示例

#### 启动 Flink SQL Client

```bash
cd $FLINK_HOME

./bin/sql-client.sh

```

#### 设置为 Batch 模式

```sql
-- 可选：建议设置运行模式为 BATCH，可获得更好的性能
SET 'execution.runtime-mode' = 'BATCH';

-- 可选：设置并行度（根据数据量和集群资源调整并行度）
SET 'parallelism.default' = '4';

```

#### 创建 DirectLoad Sink 表

在 Flink SQL 中创建对应的目标表，映射到 OceanBase 的 t_user 表：

```sql
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-directload
 - `host` 和 `port`：OceanBase 的主机地址和 RPC 端口号。连接 OBServer 服务端时，sys 租户下查询系统视图 `DBA_OB_SERVERS` 即可获取 OBServer 的 RPC 端口号。
 - `tenant-name`：租户名称
 - `username` 和 `password`：数据库用户凭证（注意：用户名不包含 @tenant 后缀）
 - `schema-name` 和 `table-name`：目标数据库和表名

#### 说明

以上使用了最小参数集。更多参数配置请参考后续章节中的 -配置参数详解。

#### 插入测试数据

```sql
-- 插入几条测试数据
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 将提交作业并开始导入数据，您会在控制台看到类似如下的输出：

```shell
[INFO] Submitting SQL update statement to the cluster...
[INFO] SQL update statement has been successfully submitted to the cluster:
Job ID: xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx

```

作业执行说明：

1. 作业将在后台执行，您可以在 Flink Web UI（http://localhost:8081）中查看作业进度。
 2. 作业状态会经历：`INITIALIZING` → `RUNNING` → `FINISHED`。
 3. DirectLoad 的最终提交发生在作业结束阶段，请务必等待作业状态变为 `FINISHED` 后再验证数据。

#### 验证数据写入

作业成功完成后，在 OceanBase 中查询验证：

```sql
SELECT * FROM test.t_user ORDER BY id;

```

查询结果如下：

```shell
+----+----------+------+-------+
| 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 集群：

```shell
# 停止 Flink 集群
cd $FLINK_HOME
./bin/stop-cluster.sh

```

在 SQL Client 中退出：

```shell
-- 在 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 中
 - **列存表优势**：导入完成后数据即为列存格式，查询性能最优，无需额外合并
 - **使用场景**：
     - 首次导入数据到空表
     - 目标表数据量很小，可接受重建
     - 追求最优的列存查询性能

```sql
-- full 模式示例
CREATE TABLE t_sink (...) WITH (
  'connector' = 'oceanbase-directload',
  'load-method' = 'full',
  ...
);

```

#### `inc`（普通增量导入）

- 适合向非空表追加数据
 - 会进行主键冲突检查，如果发现冲突按 `dup-action` 策略处理
 - **版本要求**：OceanBase 4.3.2 及以上版本
 - **限制**：暂不支持 `dup-action` 为 `REPLACE`
 - **列存表注意**：数据写入转储（转储暂不支持列存），导入后查询性能为行存性能，需要等待一次合并才能达到列存查询性能

```sql
-- 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` 模式相同，需要合并后才能达到列存查询性能

```sql
-- 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`（停止导入）

- 遇到主键冲突时立即停止导入，整个作业失败
 - **适用场景**：
     - 对数据一致性要求严格
     - 不允许出现主键冲突
     - 需要快速发现数据质量问题

```sql
'dup-action' = 'STOP_ON_DUP'

```

#### `REPLACE`（替换）

- 遇到主键冲突时，用新记录替换旧记录
 - **适用场景**：
     - 允许覆盖已存在的数据
     - 导入的是最新版本数据
 - **注意**：在 `load-method=inc` 时暂不支持

```sql
'dup-action' = 'REPLACE'

```

#### `IGNORE`（忽略）

- 遇到主键冲突时，保留旧记录，忽略新记录
 - **适用场景**：
     - 以旧数据为准
     - 增量导入时避免覆盖已存在的记录

```sql
'dup-action' = 'IGNORE'

```

### `parallel` 详解

`parallel` 参数是**服务端并行度**，它决定了 OceanBase 服务端使用多少 CPU 资源来处理本次导入任务。

#### 重要特性

1. **与客户端并发无关**：`parallel` 是服务端参数，与 Flink 的并行度是两个独立的概念
 2. **受租户配置限制**：服务端会根据租户 CPU 配置自动限制并行度上限，客户端设置超出范围不会报错
 3. **受分区分布影响**：实际并行度还受表的分区分布影响

#### 实际并行度计算规则

- **单节点并行度上限** = `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 倍**是合理的起点

```sql
-- 示例：4C 租户建议设置
'parallel' = '8'  -- 或 '16'

```

### `buffer-size` 调优建议

`buffer-size` 表示客户端写入缓冲区大小，单位是**行数**。累计到该行数后触发一次 flush 写入。

#### 调优建议

- **数据量大且单行小**：可适当调大（如 2048、4096），减少 flush 次数，提升吞吐
 - **单行数据大**：避免设置过大，防止内存压力，建议保持默认值或适当调小
 - **内存紧张**：减小 `buffer-size`，避免 TaskManager OOM

#### 推荐值

| 场景 | 推荐值 |
| --- | --- |
| 单行 < 1KB，内存充足 | 2048 – 4096 |
| 单行 1KB – 10KB | 1024（默认值） |
| 单行 > 10KB 或内存紧张 | 512 |

```sql
-- 示例
'buffer-size' = '2048'

```

## 使用示例

本节提供更完整和实用的使用示例，涵盖不同场景和配置。

### 示例 1：基础批量导入

这是最基础的使用示例，适合初学者快速上手。

```sql
-- 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` 模式）

适用于首次导入大量数据到空表的场景。

```sql
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` 模式）

适用于需要用新数据覆盖旧数据的场景。

```sql
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。

-- 创建 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。

-- 创建文件 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 任务。合理设置可以充分利用集群资源。

```sql
-- 根据数据量和集群资源设置并行度
SET 'parallelism.default' = '8';

```

#### 调优 `parallel` 参数（服务端并行度）

`parallel` 参数对 commit 阶段性能影响很大。

**调优策略**：

- 大数据量导入时，适当调大 `parallel` 可大幅缩短 commit 阶段耗时
 - 建议设置为租户 CPU 数的 **2–4 倍**
 - 不用担心设置过大，服务端会自动限制

```sql
-- 示例：8C 租户
'parallel' = '16'  -- 或 '32'

```

**性能对比示例**

| `parallel` | commit 阶段耗时 |
| --- | --- |
| 8（默认） | 约 10 分钟 |
| 16 | 约 5 分钟 |
| 32 | 约 3 分钟 |

#### 调优 `buffer-size`

根据数据特征调整缓冲区大小：

```sql
-- 小行数据（每行 < 1KB）
'buffer-size' = '4096'

-- 中等行数据（每行 1–10KB）
'buffer-size' = '1024'  -- 默认值

-- 大行数据（每行 > 10KB）
'buffer-size' = '512'

```

### 生产环境建议

#### 选择合适的导入窗口期

由于旁路导入期间会锁表，建议：

- 选择业务低峰期：如凌晨、周末
 - 提前通知相关方：避免影响其他业务
 - 设置合理的超时时间：防止导入任务长时间占用表

```sql
-- 设置2小时超时（根据数据量评估）
'timeout' = '2h'

```

#### 根据场景选择 `load-method`

| 场景 | 推荐 `load-method` | 原因 |
| --- | --- | --- |
| 首次导入到空表 | `full` | 性能最优，列存查询效果最佳 |
| 定期追加增量数据 | `inc` | 适合增量场景，支持冲突检查 |
| 每日全量更新维度表 | `inc_replace` | 自动覆盖旧数据，无需手动删除 |
| 历史数据回溯 | `full`（导入到临时表后切换） | 避免影响在线表 |

#### 列存表使用建议

如果目标是列存表，需要注意不同 `load-method` 的影响：

- **`full` 模式（推荐）**：

     - 数据直接写为列存格式
     - 导入完成后查询性能最优
     - 仅适合空表或可接受重建的表
 - **`inc` / `inc_replace` 模式**：

     - 数据写入转储（行存格式）
     - 导入后查询性能为行存性能
     - 需要等待一次 major compaction 后才能达到列存性能
     - 如果追求高查询性能，建议手动触发合并：

```sql
-- 在 OceanBase 中手动触发合并
ALTER SYSTEM MAJOR FREEZE;

```

## 常见问题

### Q1：导入期间其他写入操作失败怎么办？

**问题描述**：在执行旁路导入时，尝试对目标表进行 `INSERT`/`UPDATE`/`DELETE` 操作失败。

**原因**：这是旁路导入的固有特性。导入期间目标表会被锁定，仅允许 `SELECT` 操作。

**解决方案**：

1. **方案一：错峰导入**

      - 将导入任务安排在业务低峰期（如凌晨）
      - 提前与业务方沟通，暂停写入操作
 2. **方案二：使用中间表**

   ```sql
   -- 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;

   ```

      - 先导入到临时表
      - 导入完成后通过表切换或数据合并操作
 3. **方案三：使用普通 JDBC 连接器**

      - 如果无法接受锁表，改用 `flink-connector-oceanbase` 连接器。
      - 牺牲部分性能换取表的可用性。

### Q2：作业一直不结束 / 数据未提交怎么办？

**问题描述**：Flink 作业一直处于 `RUNNING` 状态，OceanBase 中查询不到数据。

**原因**：DirectLoad 连接器的最终 commit 发生在输入结束（end-of-input）阶段。如果输入是无界流，作业不会结束，数据也不会提交。

**排查步骤**：

1. 检查是否设置了 Batch 模式：

   ```
 2. 检查数据源是否为有界流：

      - 文件数据源：天然有界
      - Kafka：需设置 `scan.bounded.mode`
      - JDBC：天然有界
      - 自定义 Source：检查是否正确发送 end-of-input 信号
 3. 查看 Flink Web UI：

      - 检查作业状态
      - 查看是否有 Backpressure
      - 检查各个算子的处理进度

### Q3：如何选择 `load-method`？

**决策树**：

```text
是否为首次导入到空表？
├─ 是 → 使用 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 后才能转为列存格式。

**解决方案**：

1. **手动触发 major compaction（推荐）**：

   ```sql
   -- 在 OceanBase 中执行
   ALTER SYSTEM MAJOR FREEZE;

   -- 检查合并进度
   SELECT * FROM oceanbase.DBA_OB_MAJOR_COMPACTION;

   ```
 2. **等待自动合并**：  
    OceanBase 会定期自动触发合并，但可能需要等待较长时间（取决于配置）。

### Q5：导入速度很慢怎么优化？

**可能原因及解决方案**：

1. **Flink 并行度过低**：

   ```sql
   SET 'parallelism.default' = '16';  -- 增加并行度

   ```
 2. **服务端并行度过低**：

   ```sql
   'parallel' = '16'  -- 增加服务端并行度

   ```
 3. **buffer-size 过小**：

   ```sql
   'buffer-size' = '2048'  -- 增大缓冲区

   ```
 4. **OceanBase 资源不足**：

      - 检查 OceanBase 集群负载
      - 考虑扩容或选择低峰期导入

## 相关文档

- [OceanBase 旁路导入官方文档](https://www.oceanbase.com/docs/common-oceanbase-database-cn-1000000001428636)
 - [Flink 官方文档](https://nightlies.apache.org/flink/flink-docs-stable/)
 - [Flink 项目 GitHub 仓库](https://github.com/oceanbase/flink-connector-oceanbase)

 上一篇 下一篇 ![有帮助](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*y6ocSqN8cqsAAAAAAAAAAAAAARQnAQ)![无帮助](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*BG9IQJyLHF8AAAAAAAAAAAAAARQnAQ)![反馈](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*eTWdQKCRKHwAAAAAAAAAAAAAARQnAQ)[AI](https://www.oceanbase.com/obi) 咨询热线
