Apache Flinkは、分散型で高スループット、高性能なリアルタイムデータ処理シナリオ向けに設計されたオープンソースのストリーム処理フレームワークです。イベント駆動型アプリケーション、リアルタイム分析、データパイプライン、複雑なイベント処理など、さまざまなユースケースで広く利用されています。本記事では、Flink SQLクライアントとFlink CDCコネクタを使用して、ソーステーブルの変更をリアルタイムでキャプチャし、ユーザー次元情報と関連付け、その結果を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に含まれていることを確認するだけで済みます。ジョブ実行時に、-classpathや-cpなどのパラメータを使用して、コマンドライン引数でJARパスを指定できます。
操作手順
(オプション)Alibaba Cloud Flinkに接続します。
Alibaba Cloud Flinkを使用しており、かつOB CloudクラスタのクラウドベンダーもAlibaba Cloudである場合、Flinkはプライベートネットワークアドレスを介してOceanBaseに接続できます。操作手順は以下のとおりです:
Alibaba Cloudのプライベートネットワークを使用してOB Cloudに接続します。
詳細については、Alibaba Cloudプライベートネットワーク接続を使用したデータベース接続をご参照ください。
注意
VPCとVSwitchは、OB Cloudクラスタと同じリージョンにある必要があります。
リアルタイム計算Flink版を購入した後、Flink管理ページで、右上のプラグボタンをクリックし、プライベートネットワークアドレスとポートを入力して、探知をクリックします。 接続に成功すると、ネットワーク探知接続成功が返されます。
Binlogログサービスを有効にします。 Binlogログサービスを有効にするパス:インスタンスリスト -> テナント管理 -> Binlogサービスで、有効化をクリックします。詳細については、Binlogログサービスの有効化をご参照ください。
データを準備します。 サンプルデータには、orders(注文データ)とusers(ユーザー情報の次元テーブル)の2つのテーブルが含まれます。 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 に書き込まれます。このSQLは、Flink SQLクライアントまたはFlink SQL互換環境で実行することで、リアルタイム更新される注文詳細テーブルを構築できます。orders テーブルのデータに変更があると、MySQL CDCに基づくFlinkタスクが自動的にこれらの変更を ob_order_details テーブルに同期します。
コネクタのパラメータに関する詳細は、Flink公式ドキュメントをご参照ください。