首批通过分布式安全可靠测评,为关键业务系统打造
Apache Flink 与 OceanBase 数据库数据集成指南
更新时间:2026-04-09 14:12:04
本指南面向希望使用 Apache Flink 与 OceanBase 数据库构建高效数据管道的开发者。无论你是 Flink 新手还是资深用户,本文都将帮助你:
- 快速掌握 Flink 的核心使用方式。
- 全面了解 OceanBase 数据库提供的各类 Flink Connector。
- 根据实际业务场景,选择最合适的集成方案。
适用版本
- Flink V1.15 及以上。
- JDBC 连接器需要适配 OceanaBase V3.x、V4.x 及以上。
- Flink Connector OceanBase Direct Load、OBKV-HBase Connector、OBKV-HBase2 Connector 需要适配 OceanaBase V4.2.x 及以上。
Flink 基础简介
什么是 Apache Flink
Apache Flink 是一个开源的流批一体计算引擎,能够高效处理实时数据流与批量数据。在数据集成场景中,Flink 提供了强大的能力,支持跨系统间的数据同步、转换与处理。
为什么推荐 Flink SQL
Flink SQL 是使用 Flink 最简单、最高效的方式。通过类 SQL 语法,你可以:
- 连接多种数据源(Kafka、MySQL、OceanBase 数据库等)。
- 执行过滤、聚合、转换等操作。
- 将结果写入目标系统。
无需编写 Java/Scala 代码,几行 SQL 即可完成复杂任务。
从 Kafka 到 OceanBase 数据库实时同步
以下 Flink SQL 脚本可在 Flink SQL Client 或 Flink 作业中直接执行,实现从 Kafka 消费订单数据并实时同步至 OceanBase 数据库。
说明
通过 CREATE TABLE 定义源与目标,再用一条 INSERT INTO ... SELECT 语句,即可实现实时、带转换的数据同步。
定义 Kafka 源表。
CREATE TABLE kafka_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3) ) WITH ( 'connector' = 'kafka', 'topic' = 'orders', 'properties.bootstrap.servers' = 'localhost:9092', 'properties.group.id' = 'flink-consumer', 'format' = 'json', 'scan.startup.mode' = 'latest-offset' );定义 OceanBase 数据库目标表。
CREATE TABLE oceanbase_orders ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'oceanbase', 'url' = 'jdbc:mysql://127.0.0.1:2881/test', 'username' = 'root@test', 'password' = 'password', 'table-name' = 'orders' );执行数据流转(含转换与过滤)。
INSERT INTO oceanbase_orders SELECT order_id, user_id, amount * 1.1 AS amount, -- 可以进行数据转换 order_time FROM kafka_orders WHERE amount > 100; -- 可以进行数据过滤
核心概念速览
| 概念 | 说明 |
|---|---|
| Source | 数据来源(如 Kafka、OceanBase CDC)。 |
| Sink | 数据目的地(如 OceanBase 数据库、Kafka)。 |
| Connector | 连接 Flink 与外部系统的桥梁 。 |
| SQL DDL | 通过 CREATE TABLE 声明数据源/目标。 |
| SQL DML | 通过 INSERT INTO ... SELECT 驱动数据流转。 |
如何运行 Flink SQL
方式 1:SQL Client(交互式)
# 启动 Flink SQL Client ./bin/sql-client.sh # 在 Flink SQL> 提示符下执行命令 Flink SQL> CREATE TABLE ... Flink SQL> INSERT INTO ...方式 2:提交 SQL 文件
# 将 SQL 保存到文件(如 job.sql),然后提交 ./bin/sql-client.sh -f job.sql方式 3:Web UI 或编程 API
通过 Flink Web UI 提交 SQL 任务,或使用 Table API/DataStream API 编写 Java/Scala 程序。
OceanBase Flink Connector 全景图与选型指南
在上述示例中,使用了 'connector' = 'oceanbase'。实际上,OceanBase 为 Flink 提供了多种专用 Connector,覆盖读、写、CDC 等全场景。
Connector 一览表
注意
- 不同 Connector 不可互换,需根据场景精准选择。
- OceanBase 数据库 MySQL 租户兼容性较好,可以复用 Flink MySQL Connector;Oracle 租户则必须使用 Flink OceanBase Connector。
| 场景 | 推荐 Connector | 核心特性 | 文档链接 |
|---|---|---|---|
| 实时流式写入,数据量适中 | Flink Connector OceanBase | 基于 JDBC,通用性强 | Flink Connector OceanBase 官方文档 |
| Lookup 维度表关联 | Flink Connector JDBC (Lookup 模式) | 标准 JDBC | Flink Connector JDBC 官方文档 |
| 批量读取全表 | Flink Connector JDBC (批量读取,单并行度) | 标准 JDBC | Flink Connector JDBC 官方文档 |
CDC 数据同步
注意仅支持 OceanBase 数据库 MySQL 模式租户,Oracle 模式租户暂不支持。 |
Flink CDC (OceanBase CDC) | 并行全量 + 增量读取 | OceanBase CDC 官方文档 |
| 大批量数据导入,TB 级批量数据迁移 | Flink Connector OceanBase Direct Load | 基于旁路导入,高吞吐 | Flink Connector OceanBase Direct Load |
| 固定列 KV 高性能写入(简单) | Flink Connector OBKV HBase | 基于 OBKV API,嵌套结构 | Flink Connector OBKV HBase |
| 高性能 KV 写入(高级特性) | Flink Connector OBKV HBase2 | 扁平结构,支持动态列/部分更新 | Flink Connector OBKV HBase2 |
选型决策流程
可以根据下面的决策流程图选择合适的 Connector:

典型场景详解
场景 1:实时流式写入 OceanBase 数据库
需求:
将 Kafka/Pulsar 的实时数据写入 OceanBase 数据库。
方案:
Flink Connector OceanBase。详细信息,参见 Flink Connector OceanBase。
优势:
- 支持无界流(Unbounded Stream)。
- 兼容 MySQL/Oracle 模式。
- 支持批量写入与缓冲优化。
示例如下:
创建 OceanBase Sink 表。
CREATE TABLE orders_sink ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'oceanbase', 'url' = 'jdbc:mysql://127.0.0.1:2881/test', 'username' = 'root@test', 'password' = 'password', 'table-name' = 'orders', 'buffer-flush.interval' = '1s', 'buffer-flush.buffer-size' = '1000' );导入数据。
INSERT INTO orders_sink SELECT * FROM kafka_source;
场景 2:大批量数据迁移
需求:
TB 级历史数据迁移到 OceanBase 数据库。
方案:
Flink Connector OceanBase Direct Load。详细信息,参见 Flink Connector OceanBase Direct Load。
优势:
- 基于旁路导入,吞吐量极高。
- 多节点并行写入。
- 适合 Batch 模式。
注意事项:
- 支持有界流(Bounded Stream),不支持实时流。
- 导入期间目标表被锁定(只读)。
- 推荐使用 Flink Batch 模式,获取更好的性能。
示例如下:
创建 Direct Load Sink 表。
CREATE TABLE large_table_sink ( id BIGINT, name STRING, data STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'oceanbase-directload', 'url' = 'jdbc:mysql://127.0.0.1:2881/test', 'username' = 'root@test', 'password' = 'password', 'schema-name' = 'test', 'table-name' = 'large_table', 'parallel' = '8' -- 并行度 );源表批量写入结果表。
INSERT INTO large_table_sink SELECT * FROM source_table;
场景 3:高性能 KV 写入(简单场景)
需求:
需要高性能写入 KV 数据,列结构简单固定
方案:
Flink Connector OBKV HBase。详细信息,参见 Flink Connector OBKV HBase。
优势:
- 基于 OBKV HBase API,性能优异。
- 适合固定列结构的场景。
局限性:
- 表定义需要使用嵌套 ROW 结构。
- 不支持动态列。
- 不支持部分列更新。
示例如下:
创建 HBase Sink 表。
CREATE TABLE hbase_sink (
rowkey STRING,
family1 ROW<column1 STRING, column2 STRING>, -- 嵌套 ROW 结构
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'obkv-hbase',
'url' = 'http://127.0.0.1:8080/services?Action=ObRootServiceInfo&ObCluster=obcluster',
'username' = 'root@test#obcluster',
'password' = 'password',
'sys.username' = 'root',
'sys.password' = 'password',
'schema-name' = 'test',
'table-name' = 'htable1'
);
场景 4:高性能 KV 写入(高级功能)
需求:
需要高性能写入,并且需要动态列、部分列更新、灵活时间戳控制等高级特性。
方案:
Flink Connector OBKV HBase2。详细信息,参见 Flink Connector OBKV HBase2。
优势:
- 扁平化表结构,定义简洁。
- 支持动态列模式:列名可以在运行时动态指定。
- 支持部分列更新:只定义需要更新的列,未定义的列不会被更新,非常灵活。
- 支持时间戳控制:可为不同列设置不同时间戳(
tsColumn、tsMap)。 - 性能与 OBKV HBase 相当。
示例如下:
创建 HBase2 Sink 表。
基本使用:扁平结构。
CREATE TABLE hbase2_sink ( rowkey STRING, column1 STRING, -- 扁平结构,无需 ROW 嵌套 column2 STRING, PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector' = 'obkv-hbase2', 'url' = 'http://127.0.0.1:8080/services?Action=ObRootServiceInfo&ObCluster=obcluster', 'username' = 'root@test#obcluster', 'password' = 'password', 'sys.username' = 'root', 'sys.password' = 'password', 'schema-name' = 'test', 'table-name' = 'htable1', 'columnFamily' = 'f' );部分列更新示例。
假设 OceanBase 数据库中的 HBase 表有 column1, column2, column3, column4 等多个列。只需在 Flink 表中定义你想更新的列即可
CREATE TABLE partial_update_sink ( rowkey STRING, column1 STRING, -- 只定义需要更新的列 column2 STRING, -- 未在此定义的列(如 column3, column4)不会被更新 PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector' = 'obkv-hbase2', 'url' = 'http://127.0.0.1:8080/services?Action=ObRootServiceInfo&ObCluster=obcluster', 'username' = 'root@test#obcluster', 'password' = 'password', 'sys.username' = 'root', 'sys.password' = 'password', 'schema-name' = 'test', 'table-name' = 'htable1', 'columnFamily' = 'f' );
写入数据,只更新 column1 和 column2,其他列(column3, column4 等)保持不变。
INSERT INTO partial_update_sink VALUES ('1', 'new_value1', 'new_value2');
场景 5:CDC 数据同步(MySQL 租户)
需求:
实时捕获 OceanBase 数据库数据变更,实现全量+增量同步
方案:
Flink CDC (OceanBase CDC)。详细信息,参见 OceanBase CDC 官方文档。
优势:
- 并行全量读取(性能远超 JDBC)。
- 基于 binlog 的增量同步。
- 一体化全量 + 增量流程。
使用限制:
仅支持 MySQL 模式租户。
示例如下:
创建 OceanBase CDC Source 表。
CREATE TABLE orders_cdc ( order_id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), order_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'oceanbase-cdc', 'hostname' = '127.0.0.1', 'port' = '2881', 'username' = 'root@test', 'password' = 'password', 'database-name' = 'mydb', 'table-name' = 'orders', 'scan.startup.mode' = 'initial' -- 全量+增量 );读取并处理 CDC 数据。
SELECT * FROM orders_cdc;
场景 6:Lookup 维度表关联
需求:
流计算中补全用户维度信息。
方案:
Flink Connector JDBC (作为 Lookup Source)。详细信息,参见 Flink JDBC Connector 官方文档。
特点:
- 支持 Lookup Join
- 支持缓存优化
注意
全表扫描是单并行度,但 Lookup 场景通常是点查,不受此限制。
示例如下:
创建 JDBC Lookup 表。
CREATE TABLE dim_user ( user_id BIGINT, user_name STRING, city STRING, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://127.0.0.1:2881/test', 'username' = 'root@test', 'password' = 'password', 'table-name' = 'dim_user', 'lookup.cache.max-rows' = '10000', -- 缓存配置 'lookup.cache.ttl' = '1 hour' );流表关联维度表。
SELECT o.order_id, o.user_id, u.user_name, -- 从维度表补全 u.city, o.amount FROM orders_stream o LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u ON o.user_id = u.user_id;
深度对比分析
下文我们将深入分析几组容易混淆的 Connector。
OBKV HBase 与 OBKV HBase2
| 对比项 | OBKV HBase | OBKV HBase2 |
|---|---|---|
| 数据模型 | 嵌套 ROW 结构 | 扁平化模型 |
| 表定义复杂度 | 较复杂(需要 ROW<...> 嵌套) |
简洁(直接定义列) |
| 列族支持 | 一个表可写多个列族 | 一个表仅支持一个列族(多列族需为每个列族单独建表,较麻烦) |
| 动态列支持 | 不支持 | 支持 |
| 部分列更新 | 不支持 | 支持(只定义需要更新的列) |
| 时间戳控制 | 不支持 | 支持(tsColumn + tsMap) |
| 使用灵活性 | 较低 | 高 |
| 性能 | 高 | 高 |
| 学习成本 | 低(如果熟悉 HBase) | 中(需要理解新特性) |
| 适用场景 | 需要多列族的场景,简单固定列 | 单列族场景,需要动态列、部分更新、时间戳控制 |
表定义对比示例:
OBKV HBase:嵌套 ROW 结构,一个表支持多个列族。
CREATE TABLE hbase_sink ( rowkey STRING, family1 ROW<column1 STRING, column2 STRING>, -- 列族1 family2 ROW<column3 STRING, column4 STRING>, -- 列族2(支持多个列族) PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector' = 'obkv-hbase', ... );写入时可以同时写多个列族。
INSERT INTO hbase_sink VALUES ('row1', ROW('val1', 'val2'), ROW('val3', 'val4'));OBKV HBase2:扁平结构,一个表只能指定一个列族。
如果需要写多个列族,需要为每个列族创建单独的表(比较麻烦)。
列族 family1 的表:
CREATE TABLE hbase2_family1_sink ( rowkey STRING, column1 STRING, column2 STRING, PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector' = 'obkv-hbase2', 'columnFamily' = 'family1', -- 只能指定一个列族 ... );列族 family2 的表(需要单独创建)
CREATE TABLE hbase2_family2_sink ( rowkey STRING, column3 STRING, column4 STRING, PRIMARY KEY (rowkey) NOT ENFORCED ) WITH ( 'connector' = 'obkv-hbase2', 'columnFamily' = 'family2', -- 不同的列族 ... );
写入时需要分别写入两个表。
写入列族 family1 的表:
INSERT INTO hbase2_family1_sink VALUES ('row1', 'val1', 'val2');写入列族 family2 的表:
INSERT INTO hbase2_family2_sink VALUES ('row1', 'val3', 'val4');
说明
选型建议如下:
- 需要多列族:强烈推荐 OBKV HBase。
- HBase:一个表搞定多个列族,简单方便。
- HBase2:需要为每个列族创建单独的表,再分别写入,比较麻烦。
- 需要动态列、部分列更新、时间戳控制等高级特性:必须选择 OBKV HBase2。
- 简单场景(固定列,单列族):两者都可以,HBase2 表定义更简洁,推荐优先使用。
- 新项目推荐:单列族场景推荐 OBKV HBase2,多列族场景推荐 OBKV HBase。
OceanBase CDC 与 JDBC 读取
| 对比项 | OceanBase CDC | JDBC Connector |
|---|---|---|
| 全量读取并行度 | 支持并行读取(多并行度) | 单并行度(慢) |
| 增量读取 | 支持 binlog 增量读取 | 不支持 |
| Lookup Join | 不支持 | 支持(主要用途) |
| 批量全表读取 | 支持(并行,快) | 支持(单线程,慢) |
| 租户支持 | 仅 MySQL 租户 | MySQL 和 Oracle 租户 |
| 主要用途 | CDC 全量+增量同步 | Lookup 维度表关联 |
| 数据实时性 | 准实时(binlog 延迟) | 查询时实时 |
说明
选型建议如下:
- CDC 场景(全量+增量同步):优先选择 OceanBase CDC(仅 MySQL 租户)
- Lookup Join 场景(维度表关联):选择 JDBC Connector。
- 批量全表读取:优先 OceanBase CDC(MySQL 租户),Oracle 租户只能用 JDBC Connector(性能较慢)
写入 Sink 对比
| 特性 | JDBC | Direct Load | OBKV HBase | OBKV HBase2 |
|---|---|---|---|---|
| 数据流类型 | 无界/有界 | 仅有界 | 无界/有界 | 无界/有界 |
| 吞吐量 | 中 | 极高 | 高 | 高 |
| 延迟 | 低 | 高(批量) | 低 | 低 |
| 表锁定 | 无 | 导入期间锁表 | 无 | 无 |
| 兼容模式 | MySQL/Oracle | MySQL/Oracle | MySQL | MySQL |
| 复杂度 | 简单 | 简单 | 中 | 中 |
| 典型场景 | 实时写入 | 批量导入 | KV 高性能写入 | KV 高性能写入 + 高级特性 |
常见问题解答(FAQ)
问题 1:Direct Load Sink 和 JDBC Sink 的区别,如何选择?
主要区别:
- JDBC Sink:基于标准 JDBC 协议,适合 实时流式写入,支持无界流,无需锁表。
- Direct Load Sink:基于旁路导入 API,适合 大批量数据导入,吞吐量极高,但仅支持有界流且导入期间会锁表。
选择建议:
- 实时流式写入(如 Kafka -> OceanBase 数据库):选择使用 JDBC Sink。
- 大批量历史数据导入(如数据迁移):选择使用 Direct Load Sink。
问题 2:OBKV HBase 和 OBKV HBase2 的区别,如何选择?
主要区别:
表定义方式:HBase 使用嵌套 ROW 结构,HBase2 使用扁平结构(更简洁)。
列族支持:
- HBase:一个 Flink 表可以写多个列族,非常方便。
- HBase2:一个 Flink 表只能指定一个列族,如果需要写多个列族,需要为每个列族创建单独的表并分别写入,比较麻烦。
高级特性:HBase2 支持动态列、部分列更新、时间戳控制,HBase 不支持。
适用场景:HBase 适合需要多列族的场景,HBase2 适合需要高级特性的单列族场景。
选择建议:
需要多列族:强烈推荐 OBKV HBase(HBase2 需要为每个列族建表,很麻烦)。
需要动态列、部分列更新、时间戳控制场景:推荐 OBKV HBase2。
简单固定列、单列族场景:两者都可以,推荐 OBKV HBase2(表定义更简洁)。
新项目推荐:
- 单列族用 OBKV HBase2。
- 多列族用 OBKV HBase。
问题 3:JDBC Source 有哪些使用场景,与 CDC 有什么区别?
JDBC Source 主要有两种使用场景:
Lookup Join(维度表关联)- 推荐用途
- 在流计算中根据主键实时查询维度表。
- 支持缓存优化,性能高效。
- 这是 JDBC Source 的主要用途。
批量读取全表。
- 可以读取 OceanBase 数据库全表数据。
- 单并行度,性能较慢,不推荐大表使用。
- 适合小数据量的批量读取。
JDBC Source 与 CDC 的对比:
- OceanBase CDC:支持并行全量读取 + binlog 增量读取,性能高,适合 CDC 场景(仅 MySQL 租户)
- JDBC Source:主要用于 Lookup Join,也可批量读取但性能较差(单并行度)
选择建议:
- Lookup 维度表关联:选择 JDBC Source(必选)。
- CDC 全量 + 增量读取(MySQL 租户):选择 OceanBase CDC(推荐)。
- 批量读取全表:优先选择 OceanBase CDC,如不支持(Oracle 租户)再选择用 JDBC Source。
问题 4:OceanBase CDC 为什么只支持 MySQL 租户,Oracle 租户怎么办?
OceanBase CDC 支持 MySQL 兼容的 binlog service,Oracle 租户请使用 OMS 工具实现 CDC 增量读取。
更多 Binlog 的信息,参见 Binlog 服务概述。
问题 5:如何选择合适的批次大小和并行度?
批次大小(Buffer Size):
- 小批次(100-500):适合低延迟场景,数据尽快写入。
- 中批次(1000-5000):平衡延迟和吞吐量,推荐默认值。
- 大批次(5000+):适合高吞吐量场景,但延迟会增加。
并行度(Parallelism):
- 一般设置为 CPU 核心数 或其倍数。
- 旁路导入场景可以根据租户资源设置较高并行度(如 8-16)以充分利用旁路导入能力。
- 需要考虑 OceanBase 集群的负载能力。
问题 6:不同 Connector 之间可以混用吗?
可以。在一个 Flink 任务中可以同时使用多个不同的 Connector。
常见组合:
- Flink CDC (MySQL) Source + OceanBase JDBC Sink:MySQL 至 OceanBase 数据库实时同步。
- OceanBase CDC Source + Kafka Sink:OceanBase 至 Kafka 消息队列。
- Kafka Source + JDBC Lookup + OceanBase Sink:Kafka 数据关联维度表后写入 OceanBase 数据库。