首批通过分布式安全可靠测评,为关键业务系统打造
通过 Flink 同步数据
更新时间:2026-05-14 14:28:04
Apache Flink 是一个开源的流处理框架,专为分布式、高吞吐和高性能的实时数据处理场景设计。它被广泛用于事件驱动应用程序、实时分析、数据管道和复杂事件处理等各种用例。本文介绍如何使用 Flink SQL 客户端和 Flink CDC Connector 进行实时捕获源表变更,并与用户维度信息关联,将关联结果写入到 OB Cloud 的另一数据表中。
前提条件
在同步数据前,确认以下信息:
- 您已为源 OB Cloud 数据库 MySQL 兼容模式租户和目标 MySQL 数据库创建专用于数据迁移的数据库用户,并为其赋予了相关权限。
- 您已在目标 OB Cloud 数据库中创建需要的表结构,确保它与源数据库相同。
- 您已安装 Flink。更多信息,参考 Flink 下载页面。
- 您已下载 flink-sql-connector-mysql-cdc JAR 依赖文件。更多信息,参考 Flink SQL Connector MySQL CDC。
- 您已下载 flink-connector-jdbc JAR 依赖文件。更多信息,参考 Flink SQL Connector MySQL CDC。
如果您使用的是 Flink 集群模式,你需要将 JAR 文件放入 Flink 集群的 /lib 目录下。其路径通常为 $FLINK_HOME/lib,其中 FLINK_HOME 是您的 Flink 安装目录。 如果您是 Flink 单机模式,只需确保这些 JAR 文件在您的 classpath 中。您可以在运行作业时通过命令行参数指定 JAR 路径,使用诸如 -classpath 或 -cp 的参数。
操作步骤
(可选)连接阿里云 Flink。
如果您使用的是阿里云 Flink,同时您的 OB Cloud 集群的云厂商也是阿里云,那么 Flink 可以通过私网地址连接 OceanBase。操作步骤如下:
使用阿里云私网连接 OB Cloud。
更多信息,参考 使用阿里云私网连接进行数据库连接。
注意
VPC 和 VSwitch 要和 OB Cloud 集群在同一个区域。
购买实时计算 Flink 版之后,在 Flink 管理页面,点击右上方的插头按钮,输入私网地址和端口,点击 探测。 如果连接成功,将返回 网络探测连接成功。
开启 Binlog 日志服务。 开启 Binlog 日志服务的路径:实例列表 -> 租户管理 -> Binlog 服务,点击开通即可。更多信息,参考 开通 Binlog 日志服务。
准备数据。 示例数据包含 orders(订单数据)和 users(用户信息维表)两张表。 在 OB Cloud 创建表并插入数据:
-- 创建 orders 表 CREATE TABLE `orders` ( `order_id` INT PRIMARY KEY, `price` DECIMAL(10,2), `currency` VARCHAR(3), `user_id` INT ); -- 创建 users 表 CREATE TABLE `users` ( `user_id` INT PRIMARY KEY, `user_name` VARCHAR(255), `email` VARCHAR(255) ); -- 创建 order_details 表 CREATE TABLE `order_details` ( `order_id` INT PRIMARY KEY, `order_price` DECIMAL(10, 2), `currency` VARCHAR(255), `user_id` INT, `user_name` VARCHAR(255), `email` VARCHAR(255) ); -- 插入 orders 表样本数据 INSERT INTO `orders` VALUES (1, 100.00, 'USD', 1); INSERT INTO `orders` VALUES (2, 55.50, 'EUR', 2); INSERT INTO `orders` VALUES (3, 80.99, 'USD', 3); -- 插入 users 表样本数据 INSERT INTO `users` VALUES (1, 'John Doe', 'john.doe@example.com'); INSERT INTO `users` VALUES (2, 'Jane Smith', 'jane.smith@example.com'); INSERT INTO `users` VALUES (3, 'Alice Johnson', 'alice.johnson@example.com');在 Flink SQL 中,定义源表和维表以及下游的目标表。
-- 定义源表(orders) CREATE TABLE ob_orders ( order_id INT, price DECIMAL(10, 2), currency STRING, user_id INT, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 't5********.aws-ap-southeast-1.oceanbase.cloud', 'port' = '3306', 'username' = 'test', 'password' = 'xxxx', 'database-name' = 'test2', 'table-name' = 'orders' ); -- 定义维表(users) CREATE TABLE ob_users ( user_id INT, user_name STRING, email STRING, PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://t5********.aws-ap-southeast-1.oceanbase.cloud:3306/test2', 'table-name' = 'users', 'username' = 'test', 'password' = 'xxxx' ); -- 定义目标表(order_details) CREATE TABLE ob_order_details ( order_id INT, order_price DECIMAL(10, 2), currency STRING, user_id INT, user_name STRING, email STRING, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://t5********.aws-ap-southeast-1.oceanbase.cloud:3306/test2', 'table-name' = 'order_details', 'username' = 'test', 'password' = 'xxxx' );执行联合查询并将结果写入目标表。
INSERT INTO ob_order_details SELECT o.order_id, o.price, o.currency, u.user_id, u.user_name, u.email FROM ob_orders AS o JOIN ob_users AS u ON o.user_id = u.user_id;
在本查询中,通过执行 ob_orders 与 ob_users 表的 INNER JOIN 操作,依据 user_id 关联订单与用户数据。结果集包含订单的信息及用户信息,并将结果写入目标表 ob_order_details。您可在 Flink SQL 客户端或兼容 Flink SQL 的环境中运行此 SQL,以构建实时更新的订单详情表。随着 orders 表的数据变动,基于 MySQL CDC 的 Flink 任务将自动同步这些更改至 ob_order_details 表。
更多关于 Connector 的参数信息,参考 Flink 官方文档。