How Flink CDC simplifies data entry into lakes and warehouses

1. Flink CDC の概要

広義には、データの変更をキャプチャできる技術はすべて CDC 技術と呼ばれます。
一般的に CDC 技術とは、データベース内のデータ変更をキャプチャするために使用される技術を指します。
CDC 技術の適用シナリオも非常に幅広く、以下が含まれます。


*データ分散。
1つのデータソースを複数の下流に分散し、ビジネスのデカップリングやマイクロサービスによく使用されます。

*データ統合。
散在する異種データソースをデータウェアハウスに統合し、データの孤立を解消して後続の分析を容易にします。

*データ移行。
データベースのバックアップやディザスタリカバリなどによく使用されます。


Flink CDC のデータベースログに基づく Change Data Capture 技術は、全量と増分の統合読み取り機能を実現します。
Flink の優れたパイプライン機能と豊富なアップストリーム・ダウンストリームのエコシステムを活用して、さまざまなデータベースの変更をキャプチャし、それらの変更をリアルタイムでダウンストリームのストレージに同期できます。


現在、Flink CDC のアップストリームは、MySQL、MariaDB、PostgreSQL、MongoDB など豊富なデータソースをサポートしており、OceanBase、TiDB、SQL Server などのデータベースのサポートもコミュニティで計画中です。


Flink CDC のダウンストリームはさらに豊富です。
Kafka や Pulsar などのメッセージキューへの書き込みをサポートし、Hudi や Iceberg などのデータレイクへの書き込みにも対応しています。
また、さまざまなデータウェアハウスへの書き込みもサポートしています。


同時に、Flink SQL がネイティブにサポートする Changelog メカニズムを通じて、CDC データの処理を非常にシンプルにできます。
ユーザーは SQL を使ってデータベース内の全量データと増分データのクリーニング、ワイドニング、集約などの操作を実現でき、ユーザーのハードルを大幅に下げます。
さらに、Flink DataStream API はユーザーがコードを記述してカスタムロジックを実装できるようにし、サービスを深くカスタマイズする自由度を提供します。


Flink CDC 技術の核心は、テーブル内の全量データと増分データのリアルタイムで一貫性のある同期と処理をサポートし、ユーザーが各テーブルのリアルタイム一貫スナップショットを簡単に取得できるようにすることです。
たとえば、テーブルに履歴ビジネスデータの全量があり、増分のビジネスデータが継続的に書き込まれて更新されているとします。
Flink CDC は増分の更新レコードをリアルタイムでキャプチャし、データベースと一貫したスナップショットをリアルタイムで提供します。
更新レコードの場合は既存データを更新し、レコードが挿入された場合は既存データに追加されます。
このプロセス全体で、Flink CDC は一貫性を保証します。
つまり、重複や損失はありません。


では、Flink CDC 技術は既存のデータストレージおよびデータレイクアーキテクチャにどのような変化をもたらすのでしょうか。
従来のデータウェアハウスのアーキテクチャを見てみましょう。


初期のデータウェアハウスアーキテクチャでは、通常、毎日の全量データをデータウェアハウスに SELECT してからオフライン分析を行っていました。
このアーキテクチャにはいくつかの明らかな欠点があります。


毎日ビジネステーブルの全量をクエリするため、ビジネス自体の安定性に影響を与えます。

日次のオフラインスケジューリング方式では、日次出力のリアルタイム性が低下します。

クエリ方式に基づく場合、データ量が増え続けるにつれてデータベースへの負荷も継続的に増加し、アーキテクチャのパフォーマンスボトルネックが顕著です。


データウェアハウス 2.0 の時代には、データウェアハウスは Lambda アーキテクチャに進化し、増分データのリアルタイム同期インポートのリンクが追加されました。
全体的に、Lambda アーキテクチャはスケーラビリティが向上し、ビジネスの安定性に影響を与えなくなりましたが、まだいくつかの問題があります。


オフラインの定期マージに依存しており、時間単位の出力しか実現できず、遅延がまだ比較的大きい。

全量と増分が完全に別々のリンクである。

アーキテクチャ全体がリンクが長く、維持が必要なコンポーネントが多い。
このアーキテクチャのフルリンクでは DataX や Sqoop のコンポーネントを維持する必要があり、増分リンクでは Canal と Kafka のコンポーネントを維持する必要があります。
同時に、全量と増分の定期マージリンクも維持する必要があります。


従来のデータウェアハウスアーキテクチャに存在する問題に対して、Flink CDC の登場はデータウェアハウスアーキテクチャに新しい考え方をもたらしました。
Flink CDC 技術の全量増分統合リアルタイム同期機能と、データレイクが提供する更新機能を組み合わせることで、アーキテクチャ全体が非常にシンプルになります。
Flink CDC を直接使用して MySQL の全量および増分データを読み取り、Hudi に直接書き込んで更新できます。


このシンプルなアーキテクチャには明確な利点があります。
まず、ビジネスの安定性に影響を与えません。
次に、リアルタイムに近い分単位の出力を提供し、ニアリアルタイムのビジネスニーズを満たします。
同時に、全量と増分のリンクが統合され、統合同期が実現されました。
最後に、アーキテクチャのリンクが短くなり、維持が必要なコンポーネントも少なくなりました。


2. Flink CDC の中核機能

Flink CDC の中核機能は 4つの部分に分けられます。


1つ目は、増分スナップショット読み取りアルゴリズムを通じて、ロックフリー読み取り、並列読み取り、ブレークポイントからの再開などの機能を実現していることです。

2つ目は、データレイクに優しい設計であり、CDC データのレイク投入の安定性を向上させていることです。

3つ目は、異種データソースの統合をサポートし、ストリーミング ETL の処理を容易に行えることです。

4つ目は、シャーディングされたデータベースとテーブルのマージによるレイク投入をサポートすることです。
次に、これらの機能をそれぞれ紹介します。


Flink CDC 1.x バージョンでは、MySQL CDC には 3つの大きな課題があり、本番環境での可用性に影響を与えていました。


1つ目は、MySQL CDC が全量データと増分データの一貫性を保証するためにグローバルロックを使用する必要があり、MySQL のグローバルロックはオンラインビジネスに影響を与えることです。

2つ目は、単一の並列読み取りしかサポートされておらず、大規模なテーブルの読み取りに非常に時間がかかることです。

3つ目は、全量同期ステージで、ジョブ障害発生時に再同期しかできず、安定性が悪いことです。
これらの問題に対応するため、Flink CDC コミュニティは「増分スナップショット読み取りアルゴリズム」を提案し、同時にロックフリー読み取り、並列読み取り、ブレークポイントからの再開の機能を実現し、これらの課題をまとめて解決しました。


簡単に言うと、増分スナップショット読み取りアルゴリズムの核心的な考え方は、全量読み取りフェーズでテーブルをチャンクに分割して並列読み取りを行い、増分フェーズに移行後は 1つのタスクで Binlog ログを並列に読み取るだけで、全量から増分への自動切り替え時にロックフリーアルゴリズムを通じて一貫性を保証するというものです。
この設計は、読み取り効率の向上と同時にリソースをさらに節約し、データ同期の全量増分統合を実現しました。
これはストリームバッチ統合の道筋において非常に重要な成果です。


Flink CDC はストリーミングレイク対応フレームワークです。
Flink CDC の初期設計バージョンでは、データレイクシナリオは考慮されておらず、全量フェーズでチェックポイントがサポートされていなかったため、全量データが 1つのチェックポイントで処理され、チェックポイントに依存してデータを送信するデータレイクにとって非常に不利でした。
Flink CDC 2.0 は、設計段階からデータレイクシナリオを念頭に置いて設計されており、ストリームからレイクへの対応を考慮した設計です。
設計上、全量データは細かく分割され、Flink CDC はチェックポイントの粒度をテーブル粒度からチャンク粒度に最適化できます。
これにより、データレイク書き込み時のバッファ使用量が大幅に削減され、データレイクへの書き込みにとってより友好的です。


Flink CDC を他のデータ統合フレームワークと差別化する核心的なポイントの 1つは、Flink が提供するストリームバッチ統合コンピューティング機能です。
これにより、Flink CDC は完全な ETL ツールとなり、優れた E および L の機能だけでなく、強力な Transformation 機能も備えています。
したがって、異種データソースに基づくデータレイクの構築を簡単に実現できます。


左側の SQL では、MySQL のリアルタイムプロダクトテーブルとリアルタイム注文テーブルを PostgreSQL のリアルタイム物流情報テーブルと関連付け、つまりストリーミング Join を実行しています。
関連付け結果はリアルタイムで Hudi に更新され、異種データソースのデータレイク構築を非常に簡単に効率的に完了できます。


OLTP システムでは、1つのテーブルの大量のデータの問題を解決するために、通常シャーディング方式が採用され、単一の大きなテーブルを分割してシステムのスループットを向上させます。
しかし、データ分析を容易にするため、通常、データウェアハウスやデータレイクに同期する際に、シャーディングされたテーブルを 1つの大きなテーブルにマージする必要があります。
Flink CDC はこのタスクを簡単に実行できます。


左側の SQL では、すべてのユーザーサブデータベースとサブテーブルのデータをキャプチャする user_source テーブルを宣言しています。
テーブル設定項目 database-name と table-name を通じて正規表現を使用してこれらのテーブルをマッチングしています。
さらに、user_source テーブルにはデータの出所を区別するための 2つのメタデータ列が定義されています。
Hudi テーブルの宣言では、ライブラリ名、テーブル名、および元のテーブルのプライマリキーを Hudi 内の複合プライマリキーとして宣言しています。
2つのテーブルを宣言した後、シンプルな INSERT INTO 文で、すべてのサブデータベースとサブテーブルのデータを Hudi 内の 1つのテーブルにマージでき、シャーディングに基づくデータレイクの構築を完了します。
これにより、後続のレイク上での統合分析が容易になります。


3. Flink CDC のオープンソースエコロジー

Flink CDC は独立したオープンソースプロジェクトで、プロジェクトコードは GitHub でホストされています。
小刻みに素早く進めるリリースリズムを採用し、コミュニティは今年 5つのバージョンをリリースしました。
1.x シリーズの 3つのバージョンではいくつかの小機能が導入されました。
2.0 バージョンでは MySQL CDC がロックフリー読み取り、並列読み取り、ブレークポイントからの再開などの高度な機能をサポートし、91 コミット、15 人の貢献者が参加しました。
バージョン 2.1 では MongoDB データベースをサポートし、115 コミット、28 人の貢献者が参加しました。
コミュニティのコミット数と貢献者数は大幅に増加しました。


ドキュメントやヘルプマニュアルもオープンソースコミュニティの非常に重要な部分です。
ユーザーをより良くサポートするため、Flink CDC コミュニティはバージョン管理されたドキュメントウェブサイトを立ち上げました。
たとえば、バージョン 2.1 のドキュメントがあります。
ドキュメントにはクイックスタートチュートリアルも多数提供されており、Docker 環境があれば Flink CDC をすぐに使い始められます。
さらに、よくある質問のマニュアルが提供されており、ユーザーが遭遇する一般的な問題を迅速に解決できます。


2021 年、Flink CDC コミュニティは急速な発展を遂げ、GitHub の PR と Issue は非常に活発で、GitHub Star は前年比 330% 増加しました。


4. Alibaba における Flink CDC の実践と改善

Flink CDC のレイク投入とデータウェアハウスへの適用は、Alibaba でも大規模に実践・実装されており、プロセスの中でいくつかの課題やチャレンジにも直面してきました。
どのように改善し、解決したかを紹介します。


まず、CDC がレイク投入時に直面した課題とチャレンジを見てみましょう。
これはユーザーの元の CDC データレイクアーキテクチャで、2つのリンクに分かれています。


全量データを一度に取得する全量同期ジョブがあります。

Binlog データを Canal と処理エンジンを通じて Hudi テーブルに準リアルタイムで同期する増分ジョブもあります。

このアーキテクチャは Hudi の更新機能を活用しているため、定期的にフルマージタスクをスケジュールする必要がなく、分単位の遅延を実現できます。
しかし、全量と増分は依然として 2つの別々の操作であり、全量と増分の切り替えには引き続き手動介入が必要で、正確な増分開始ポイントを指定する必要があります。
そうでなければ、データ損失のリスクがあります。
このアーキテクチャはストリームとバッチに分かれており、統一された全体ではないことがわかります。
Xuejin も紹介したように、Flink CDC の最大の利点の 1つは全量増分の自動切り替えであり、ユーザーの元のレイク内アーキテクチャを Flink CDC に置き換えました。


しかし、ユーザーが Flink CDC を使用した後、最初に直面する課題は、MySQL DDL を Flink DDL に手動でマッピングする必要があることです。
テーブル構造の手動マッピングは煩雑で、特にテーブル数やフィールド数が非常に多い場合に負担が大きいです。
さらに、手動マッピングはエラーも発生しやすくなります。
たとえば、MySQL の BIGINT UNSIGNED は Flink の BIGINT にマッピングできず、DECIMAL(20) にマッピングする必要があります。
システムがユーザーのテーブル構造を自動的にマッピングできれば、はるかにシンプルで安全になります。


ユーザーが直面するもう 1つの課題は、テーブル構造の変更により、レイクへのリンクの維持が難しくなることです。
たとえば、ユーザーに元々 id と name の 2列しかなかったテーブルがあるとします。
しかし、突然 Address 列が追加されました。
新しく追加された列のデータはデータレイクに同期されない可能性があり、レイクへのリンクが停止して安定性に影響を与える可能性もあります。
列の追加変更に加えて、列の削除や型の変更なども考えられます。
Fivetran は海外で調査レポートを発表し、企業の 60% は毎月スキーマが変更され、30% は毎週変更されることを発見しました。
これは、基本的にすべての企業がスキーマ変更に伴うデータ統合の課題に直面していることを示しています。


最後は、データベース全体をレイクに投入する際の課題です。
ユーザーは主に SQL を使用しているため、各テーブルのデータ同期リンクに対して INSERT INTO 文を定義する必要があります。
ユーザーの中には、MySQL インスタンスに数千のビジネステーブルを持つ場合もあり、数千の INSERT INTO 文を書かなければなりません。
さらに厄介なのは、各 INSERT INTO タスクは少なくとも 1つのデータベース接続を確立し、Binlog データを 1回読み取るということです。
数千のテーブルがレイクに投入される場合、数千の接続と数千回の繰り返しの Binlog 読み取りが必要になります。
これは MySQL とネットワークに大きな負荷をかけます。


これまで CDC データのレイク投入に関する多くの課題とチャレンジを紹介してきました。
ユーザーの視点から考えてみましょう。
データベースをレイクに投入するシナリオで、ユーザーは実際に何を求めているのでしょうか。
中間のデータ統合システムをブラックボックスと見なすことができます。
ユーザーはこのブラックボックスにどのような機能を提供してほしいと期待しているのでしょうか。レイクへの投入を簡素化するためです。

まず、ユーザーは間違いなくデータベース内の全量データと増分データの両方を同期したいと考えており、これにはシステムに全量増分の統合と全量増分の自動切り替え機能が必要です。分割された全量リンクと増分リンクではありません。

次に、ユーザーは各テーブルのスキーマを手動でマッピングしたくないと考えており、これにはシステムにメタ情報を自動的に検出する機能が必要です。ユーザーが Flink で DDL を作成する手間を省き、さらにユーザーの代わりに Hudi にターゲットテーブルを自動的に作成することが求められます。

さらに、ユーザーはソーステーブル構造の変更も自動的に同期されることを望んでいます。列の追加、列の削除、列の変更であっても、テーブルの追加、テーブルの削除、テーブルの変更であっても、すべてリアルタイムで自動的にターゲット側に同期され、ソース側で発生する新しいデータを失うことなく、ソース側のデータベースとデータ整合性を維持する ODS レイヤーを自動的に構築します。

最後に、本番環境で利用可能なデータベース全体の同期機能が必要です。ソース側に過度な負荷をかけてオンラインビジネスに影響を与えないようにするためです。

これら 4つの中核機能は、ユーザーが期待する理想的なデータ統合システムを基本的に構成しており、これらすべてを 1行の SQL と 1つのジョブで完了できれば、さらに完璧です。中間のこのシステムを「完全自動データ統合」と呼んでいます。データベースのレイク投入を自動的に完了し、現在直面しているいくつかの中核的な課題を解決するからです。そして、Flink はこの目標を達成するために非常に適したエンジンのようです。

そこで、Flink に基づくこの「完全自動データ統合」の構築に多くの力を注ぎました。主に先ほど述べた 4つのポイントを中心に展開しています。

まず、Flink CDC には既に全量増分の自動切り替え機能があり、これも Flink CDC のハイライトの 1つです。

メタ情報の自動検出については、Flink の Catalog インターフェースを通じてシームレスに接続できます。MySQL Catalog を開発して MySQL 内のテーブルとスキーマを自動的に検出し、Hudi Catalog も開発して Hudi 内のターゲットテーブルのメタ情報を自動的に作成できるようにしました。

テーブル構造変更の自動同期については、Schema Evolution カーネルを導入し、Flink ジョブが外部サービスに依存せずにスキーマ変更をリアルタイムで同期できるようにしました。
データベース全体の同期については、CDAS 構文を導入し、1行の SQL 文でデータベース全体の同期ジョブの定義を完了でき、ソースマージの最適化も導入してソースデータベースへの負荷を軽減しました。


データベース全体の同期をサポートするため、CDAS と CTAS のデータ同期構文を導入しました。構文は非常にシンプルです。CDAS 構文は create database as database で、主にデータベース全体の同期に使用されます。ここに表示されている文は、MySQL の tpc_ds ライブラリから Hudi の ods ライブラリへのデータベース全体の同期を完了します。同様に、CTAS 構文もあり、テーブルレベルの同期を簡単にサポートでき、正規表現でライブラリ名とテーブル名を指定してシャーディングのマージ同期を完了できます。ここにあるように、MySQL のユーザーサブデータベースとサブテーブルを Hudi の users テーブルにマージします。CDAS と CTAS の構文は自動的にターゲットに移動してターゲットテーブルを作成し、その後 Flink ジョブを開始して全量データと増分データを自動的に同期し、テーブル構造の変更もリアルタイムで同期します。

前述のように、数千のテーブルがレイクにインポートされる際、多数のデータベース接続が確立され、Binlog の繰り返し読み取りによってソースデータベースに大きな負荷がかかります。この問題を解決するため、ソースマージの最適化を導入しました。同じジョブ内でソースのマージを試み、同じデータソースを読み取る場合は 1つのソースノードにマージします。この時、データベースは 1つの接続を確立するだけで済み、Binlog も 1回だけ読み取られます。これにより、データベース全体の読み取りが実現され、データベースへの負荷が軽減されます。

データレイクおよびデータウェアハウスへのデータ投入の作業をどのように簡素化するかをより直感的に理解してもらうため、追加のデモビデオも提供しています。興味のある方は、Flink Forward Asia 2021 カンファレンスでの「How Flink CDC Simplifies Real-time Data Entry into the Lake Warehouse」の発表をご覧ください。

5. Flink CDC の将来計画

最後に、Flink CDC の将来の計画は主に 3つあります。

1つ目は、CDAS と CTAS の構文とインターフェースを継続的に改善し、Schema Evolution のコアを磨き上げ、オープンソースの準備を進めることです。

2つ目は、TiDB、OceanBase、SQL Server を含むより多くの CDC データソースを拡張していくことです。これらは既に計画段階にあります。

3つ目は、現在の増分スナップショット読み取りアルゴリズムを汎用フレームワークとして抽象化し、より多くのデータベースがいくつかのシンプルなインターフェースを通じてこのフレームワークに接続できるようにし、全量増分統合の機能を持たせることです。

Related Articles

Explore More Special Offers

  1. Short Message Service(SMS) & Mail Service

    50,000 email package starts as low as USD 1.99, 120 short messages start at only USD 1.00

phone お問い合わせ
Hi, I'm Alibaba Cloud AI Assistant!
I can help with questions and solutions.