データサブスクリプションは、変更データキャプチャ(CDC)を通じて、上流からデータベースログ(MySQL Binlog、OceanBase Binlogサービスなど)を継続的に読み取りし、増分変更を解析・ソートまたは変換した後、ターゲット側で再生します。一般的なターゲットには、OceanBase、Kafkaが含まれ、さらにFlinkやカスタムコンシューマーを経由して分析データベースやワイドテーブルに書き込まれます。
代表的なユースケース:BI/データウェアハウスへのリアルタイム同期、キャッシュと検索エンジンの更新、移行切り替え期間中の増分データの追いつけ、Kafkaを介したトラフィックのピークカットと複数チームによる消費、OceanBase Binlogの外部出力などです。
Canalを使用したMySQL Binlogのサブスクリプション
CanalはMySQLスレーブがBinlogをサブスクライブすることをシミュレートし、行レベルの変更をクライアントに配信します。OceanBaseへの直接同期を設定するか、TCP / Kafka / RocketMQ / RabbitMQなどのモードでメッセージシステムに先に書き込み、その後下流で消費され、ターゲットデータベース(OceanBaseを含む)に書き込まれます。
ケース |
説明 |
ドキュメント |
|---|---|---|
| MySQL → OceanBase(Canal) | MySQL Binlogの解析に基づき、OceanBaseに増分書き込みを行います。 | Canalを使用したMySQLデータベースからOceanBaseデータベースへのデータ同期 |
| OceanBase → MySQL(Canal) | ソース側にはOceanBase Binlogサービスが必要で、ログをMySQL Binlog形式に準じたものとしてサブスクリプションに供給します。 | Canalを使用したOceanBaseデータベースからMySQLデータベースへのデータ同期 |
| MySQL → OceanBase(CloudCanal) | コミュニティ版パイプラインで、Binlog系の同期ソリューションに属します。Canalと比較して選定できます。 | CloudCanalを使用したMySQLデータベースからOceanBaseデータベースへのデータ移行 |
Flink CDCを使用した複数ソースデータベースのサブスクリプション
Flink CDCは、Flink上でDatabase CDC Sourceを提供し、様々なデータベースからストアド + 増分を取得できます。Flink SQLと組み合わせて関連付け、ワイド化、集計を行った後、Kafka、JDBC(OceanBaseを含む)、HiveなどのSinkに書き込みます。マルチソース統合、リアルタイム処理、およびExactly-Once要件が高いパイプラインに適しています。
ケース |
説明 |
ドキュメント |
|---|---|---|
| MySQL → OceanBase | MySQL CDC Sourceを使用し、OceanBaseに増分同期します。 | Flink CDCを使用したMySQLデータベースからOceanBaseデータベースへのデータ同期 |
| OceanBase → MySQL | OceanBase Binlogサービスに依存し、Flink内にOceanBase CDCソーステーブルを作成します。 | Flink CDCを使用したOceanBaseデータベースからMySQLデータベースへのデータ移行 |
| OceanBase増分(ChunJun / FlinkX) | Binlogサービスとoblogclientに基づく増分同期フレームワークです。 |
ChunJunを使用したOceanBaseデータベースからMySQLデータベースへのデータ移行 |
| StarRocks → OceanBase(Flinkジョブ) | OceanBaseが提供するFlink移行ツールで、SR処理パイプラインのベーステーブルが一致しない場合に適しています。 | Flink-OMTを使用したStarRocksデータベースからOceanBaseデータベースへのデータ同期 |
データサブスクリプションを使用したKafka / OceanBaseへの書き込みシナリオ
一般的なパターンは次のとおりです:CDCまたはCanalがログを解析 → Kafkaに書き込む(バッファリング、複数サブスクリプション)→ Flink/OMS/自社開発コンシューマー → OceanBaseまたは他のストレージに落とす。また、Flink CDCが直接OceanBase JDBCにSinkすることも可能で、Kafkaを省略できます。スループットと結合度に応じて選択します。
ケース |
説明 |
ドキュメント / リンク |
|---|---|---|
| TiDB:TiCDC → Kafka → OMS → OceanBase | TiDBの増分データはKafkaを経由して配信され、OMSによってOceanBase MySQLテナントに書き込まれます(新しいデータソースを作成する場合はKafkaにバインドする必要があります)。 | OMSを使用したTiDBデータベースからOceanBaseデータベースへのデータ移行(MySQL互換モードテナント) |
| OMSがKafka / RocketMQ / Datahubに接続 | 製品側がサポートする下流メッセージとデータ統合の形態で、サブスクリプションパイプラインと組み合わせることができます。 | データ移行の概要 · OMSがサポートするデータソース(移行ソリューションのサポート状況表);OMS使用ドキュメント |
| Canal → Kafka(再消費してOBに書く) | CanalのserverModeなどがkafkaなどに設定されている場合、変更はまずTopicに入り、その後コンシューマーによってOceanBaseに書き込まれます。 |
Canalを使用したMySQLデータベースからOceanBaseデータベースへのデータ同期 |
| Flink CDC → Kafka → 下流 | FlinkはCDCストリームをKafkaに書き込み、リアルタイムデータウェアハウスや複数チームの消費に供します。Sinkの選定についてはFlinkドキュメントを参照してください。 | Flink Kafka Connector |
| Flink / JDBCによるOceanBaseへの直接書き込み | Kafkaを経由せず、Flink JDBCまたはOceanBase関連のConnectorを使用して内部テーブルに直接書き込みます(インポートやダイレクトロード機能と組み合わせる際は、意味合いに注意してください)。 | Flink JDBC SQL Connector;インポート側については データインポートの概要 を参照してください |