Practice and optimization of OLAP analysis and real-time data warehouse by BIGO using Flink

1. 業務背景

BIGO は主にショート動画ライブ配信事業を手掛ける海外志向の企業です。現在、主な事業として BigoLive(グローバルライブ配信サービス)、Likee(ショート動画作成・共有プラットフォーム)、IMO(無料コミュニケーションツール)を展開しており、中国で 4 億人のユーザーを擁しています。事業の発展に伴い、データプラットフォームの処理能力に対する要求がますます高まっており、プラットフォームが直面する課題もいっそう顕著になっています。以下では、BIGO ビッグデータプラットフォームとその直面する課題について紹介します。BIGO ビッグデータプラットフォームのデータフロー図は以下の通りです。

APP およびウェブページ上のユーザー行動ログデータ、ならびにリレーショナルデータベースの Binlog データは、BIGO ビッグデータプラットフォームのメッセージキューとオフラインストレージシステムに同期され、リアルタイムおよびオフラインのデータ分析手段を経て計算され、リアルタイムレコメンデーション、モニタリング、アドホッククエリなどのシナリオに適用されます。ただし、以下の課題が存在します。

*OLAP 分析プラットフォームの入口が統一されていない:Presto/Spark の分析タスク入口が並存しており、ユーザーは自身の SQL クエリに適したエンジンを判別できず、やみくもに選択するため、体験が良くありません。さらに、ユーザーは同じクエリを 2 つの入口から同時に送信して迅速に結果を取得しようとするため、リソースの無駄が生じています。
*オフラインタスクの計算遅延が大きく、結果の出力が遅すぎる:ABTest などの典型的なビジネスシナリオでは、結果の計算が午後までかかることがよくあります。
*各ビジネス部門が独自のシナリオに基づいてアプリケーションを個別に開発し、リアルタイムタスクがサイロ型で開発されているため、データの階層化やデータリネージが欠如しています。

上記の課題に対処するため、BIGO ビッグデータプラットフォームは OneSQL OLAP 分析プラットフォームとリアルタイムデータウェアハウスを構築しました。

*OneSQL OLAP 分析プラットフォームを通じて、OLAP クエリ入口を統一し、ユーザーの選択の迷いを減らし、プラットフォームのリソース利用率を向上させます。
*Flink でリアルタイムデータウェアハウスタスクを構築し、Kafka/Pulsar でデータ階層化を行います。
*計算が遅い一部のオフラインタスクを Flink ストリーミングコンピューティングタスクに移行し、計算結果の出力を高速化します。
*さらに、リアルタイムコンピューティングプラットフォーム Bigoflow を構築してこれらのリアルタイムコンピューティングタスクを管理し、リアルタイムタスクのデータリネージを構築します。

2. 実践と特性改善

2.1 OneSQL OLAP 分析プラットフォームの実践と最適化
OneSQL OLAP 分析プラットフォームは、Flink、Spark、Presto を統合した OLAP クエリ分析エンジンです。ユーザーが投稿した OLAP クエリリクエストは、OneSQL バックエンドを通じて異なる実行エンジンのクライアントに転送され、対応するクエリリクエストが異なるクラスターに送信されて実行されます。全体構造図は以下の通りです。

分析プラットフォームの全体構造は、上から入口層、転送層、実行層、リソース管理層に分かれています。ユーザー体験の最適化、実行失敗の確率低減、各クラスターのリソース利用率向上のため、OneSQL OLAP 分析プラットフォームは以下の機能を実装しています。

*統一クエリ入口:入口層で、ユーザーは Hive SQL 構文を標準とする統一された Hue クエリページ入口を通じてクエリを投稿します。
*統一クエリ構文:Flink、Spark、Presto など複数のクエリエンジンを統合し、各クエリエンジンが Hive SQL 構文に適応することでユーザーの SQL クエリタスクを実行します。
*インテリジェントルーティング:実行エンジンの選定過程で、過去の SQL クエリの実行状況(各エンジンでの実行成否と実行時間)、各クラスターの負荷状況、各エンジンの SQL 構文互換性に基づき、適切なエンジンを選択してクエリを投稿します。
*失敗時リトライ:OneSQL バックエンドは SQL タスクの実行を監視し、SQL タスクが実行中に失敗した場合、別のエンジンを選択してリトライを実行します。

このように、BIGO ビッグデータプラットフォームは OneSQL OLAP 分析プラットフォームを通じて、OLAP 分析ポータルの統一を実現し、ユーザーの選択の迷いを減らし、各クラスターのリソースを有効活用してリソースアイドルを削減しています。

2.1.1 Flink OLAP 分析システムの構築

OneSQL 分析プラットフォームにおいて、Flink は OLAP 分析エンジンの一部としても機能しています。Flink OLAP システムは Flink SQL Gateway と Flink Session クラスターの 2 つのコンポーネントで構成されています。SQL Gateway は SQL 投稿の入口として機能し、クエリ SQL は Gateway を通じて Session クラスターに送信されて実行され、同時に SQL 実行の進捗状況を取得してクエリ結果をクライアントに返します。SQL クエリの実行プロセスは以下の通りです。

まず、ユーザーが投稿した SQL は SQL Gateway で判定されます。結果を Hive テーブルに永続化する必要がある場合は、HiveCatalog インターフェイスを通じて Hive テーブルを作成し、クエリタスクの計算結果を永続化します。その後、タスクは SQL Gateway で SQL 解析を実行し、ジョブオペレーションの並列度を設定して、パイプラインを生成し Session クラスターに送信して実行します。

Flink OLAP システム全体の安定性を確保し、SQL クエリを効率的に実行するため、本システムでは以下の機能強化を実施しました。

安定性:

ZooKeeper HA ベースで Flink Session クラスターの信頼性を保証し、SQL Gateway が ZooKeeper ノードを監視して Session クラスターを検知します。
クエリが Hive テーブルをスキャンするデータ量、パーティション数、および戻りデータ量を制御し、Session クラスターの JobManager と TaskManager のメモリ不足を防止します。
パフォーマンス:

Flink Session クラスターがリソースを事前割り当てし、ジョブ投稿後のリソース申請時間を短縮します。
Flink JobManager がスプリットの解析を非同期で行い、タスクの解析中にスプリットを実行することで、スプリット解析待ちによるタスクの実行時間を短縮します。
ジョブ投稿過程でスキャンパーティションとスプリットの最大数を制御し、並列タスクのセットアップ時間を短縮します。
Hive SQL 互換性:
Flink の Hive SQL 構文互換性を改善し、現在 Hive SQL との互換性は約 80% に達しています。

モニタリングとアラート:
Flink Session クラスターの JobManager、TaskManager、SQL Gateway のメモリ、CPU 使用率、タスク投稿状況を監視し、問題が発生した場合は直ちにアラートで通知して対応します。

2.1.2 OneSQL OLAP 分析プラットフォームの成果
上記の実装に基づく OneSQL OLAP 分析プラットフォームは、以下のメリットを実現しました。

クエリ入口を統一し、ユーザーの選択ミスを削減。ユーザーの実行エラー率は 85.7% 減少し、SQL 実行成功率は 3% 向上しました。
SQL 実行時間は 10% 短縮され、各クラスターのリソースを有効活用してタスクのキュー待ち時間を削減しました。
OLAP 分析エンジンの一部としての Flink が、リアルタイムコンピューティングクラスターのリソース利用率を 15% 向上させました。
2.2 リアルタイムデータウェアハウスの構築と最適化
BIGO ビッグデータプラットフォームにおいて一部のビジネス指標の出力効率を改善し、Flink リアルタイムタスクをより適切に管理するため、BIGO ビッグデータプラットフォームはリアルタイムコンピューティングプラットフォーム Bigoflow を構築し、計算が遅いタスクをリアルタイムコンピューティングプラットフォームに移行しました。Flink ストリーミングコンピューティングで実行し、メッセージキュー Kafka/Pulsar を通じてデータ階層化を行い、リアルタイムデータウェアハウスを構築しています。Bigoflow 上では、リアルタイムデータウェアハウスのタスクをプラットフォームベースで管理し、統一されたリアルタイムタスクアクセス入口を設けています。さらに、プラットフォームベースでリアルタイムタスクのメタデータを管理し、リアルタイムタスクのデータリネージを構築しています。

2.2.1 構築プラン
BIGO ビッグデータプラットフォームは主に Flink + ClickHouse をベースにリアルタイムデータウェアハウスを構築しています。全体のプランは以下の通りです。

従来のデータウェアハウスのデータ階層化アプローチに従い、データを ODS、DWD、DWS、ADS の 4 層に分割します。

ODS レイヤー:ユーザー行動ログ、ビジネスログなどの生データを Kafka/Pulsar などのメッセージキューに格納します。
DWD レイヤー:Flink タスクでユーザーの UserId ごとに集約し、異なるユーザーの詳細な行動データを形成して Kafka/Pulsar に保存します。
DWS レイヤー:ユーザー行動詳細を持つ Kafka フローテーブルとユーザーの Hive/MySQL ディメンションテーブルでフローディメンションテーブル JOIN を実行し、JOIN 後に生成された多次元詳細データを ClickHouse テーブルに出力します。
ADS レイヤー:ClickHouse 内の多次元詳細データを異なるディメンションで集約し、異なるビジネスに適用します。
上記のプランに従ってリアルタイムデータウェアハウスを構築する過程で、いくつかの課題に直面しました。

オフラインタスクをリアルタイムコンピューティングタスクに変換後、計算ロジックが複雑になり(マルチストリーム JOIN、重複排除)、ジョブステートが大きすぎてメモリ不足例外が発生したり、ジョブ演算子に過度なバックプレッシャーが生じたりする問題があります。
ディメンションテーブルの JOIN 過程で、詳細フローテーブルと大規模ディメンションテーブルの結合時に、ディメンションテーブルのデータ量が大きすぎてメモリへの読み込みでメモリ不足が発生し、ジョブが失敗して実行できない問題があります。
Flink がフローディメンションテーブル JOIN で生成した多次元詳細データを ClickHouse に書き込む際、Exactly-Once を保証できず、ジョブ障害時にデータの重複書き込みが発生する問題があります。

2.2.2 課題解決と最適化

ジョブ実行ロジックの最適化とステート削減

オフラインコンピューティングタスクのロジックは複雑で、複数の Hive テーブル間の JOIN や重複排除操作が含まれます。一般的なロジックは以下の通りです。

オフラインジョブを Flink ストリーミングタスクに変換後、複数の Hive テーブルを JOIN していたオフラインシナリオは、複数の Kafka トピックを JOIN するシナリオに変更されます。JOIN 対象の Kafka トピックのトラフィックは大きく、JOIN のウィンドウ時間が長いため(最長ウィンドウは 1 日)、ジョブを一定期間実行すると、JOIN 演算子に大量のステートが蓄積されます(1 時間後にステートは約 1 TB に近づきます)。このような大規模ステートに対し、Flink ジョブは RocksDB ステートバックエンドを使用してステートデータを保存していますが、RocksDB のメモリ使用量が過大で YARN により強制終了される問題や、RocksDB に保存されたステート量が過多でスループットが低下し、ジョブに深刻なバックプレッシャーが発生する問題を避けられませんでした。

この問題を解決するため、これら複数のトピックを同じスキーマで UNION して大規模なデータストリームを生成し、このデータストリーム内の異なるイベントストリームの event_id に基づいて判定することで、各データがどのイベントストリームのトピックに由来するかを判別し、集約計算を実行して対応するイベントストリーム上の計算指標を取得します。

このように、JOIN を UNION ALL に置き換えることで、JOIN 計算による大規模ステートの影響を回避しました。

さらに、計算タスク内には以下のような複数の count distinct 計算が存在します。

これらの count distinct 計算は同じ group by 内にあり、同じ postid に基づいて重複排除計算が行われるため、これらの distinct ステートは 1 セットのキーを共有して重複排除計算を行うことができます。そこで、これらの count distinct のステートを 1 つの MapState で保存できます。

これらの count distinct 関数は同じキーに対して重複排除を行うため、MapState 内のキー値を共有してストレージスペースを最適化できます。一方、MapState の Value は Byte 配列で、各 Byte は 8 ビットを持ち、各ビットは 0 または 1 です。n 番目のビットは n 番目の count distinct 関数の対応するキーでの値に対応します。1 はその count distinct 関数が対応するキーでカウントする必要があることを示し、0 はカウント不要を示します。集約結果を計算する際、すべてのキーのビットを加算したものが n 番目の count distinct の値となり、ステートのストレージスペースをさらに節約できます。

上記の最適化により、ABTest のオフラインタスクを Flink ストリーミングコンピューティングタスクに正常に移行し、ジョブのステートを 100 GB 以内に制御して、ジョブの正常実行を実現しました。

フローディメンションテーブル JOIN の最適化

多次元詳細ワイドテーブルを生成する過程で、フローディメンションテーブルの JOIN が必要であり、Flink の Join Hive ディメンションテーブル機能を使用します。Hive ディメンションテーブルのデータはタスクの HashMap のメモリ内データ構造に読み込まれ、フローテーブルのデータは Join キーと HashMap 内のデータに基づいて JOIN されます。ただし、数千万行や数億行の Hive 大規模ディメンションテーブルの場合、メモリに読み込むデータ量が大きすぎて、メモリ不足が発生しやすくなります。上記の課題に対処するため、Join キーに基づいて Hive 大規模ディメンションテーブルをハッシュパーティショニングしました。

このように、Hive 大規模ディメンションテーブルのデータはハッシュ関数で計算され、Flink ジョブの異なる並列サブタスクの HashMap に分散されます。各 HashMap は大規模ディメンションテーブルの一部のデータのみを保存します。ジョブの並列度が十分に高ければ、大規模ディメンションテーブルのデータを十分な数の断片に分割して断片化して保存できます。一部の大規模すぎるディメンションテーブルについては、RocksDB マップステートを使用して断片化データを保存することもできます。

Kafka フローテーブルのデータを異なるサブタスクに送信して JOIN を実行する際も、同じ Join キーを通じて同じハッシュ関数で計算され、データが対応するサブタスクに割り当てられて JOIN が実行され、JOIN 後の結果が出力されます。

上記の最適化により、一部の Hive 大規模ディメンションテーブルタスクでのフローディメンションテーブル結合計算が正常に実行できるようになりました。最大のディメンションテーブルは 10 億行を超えています。

ClickHouse Sink の Exactly-Once セマンティクスサポート

フローディメンションテーブル JOIN で生成された多次元詳細データを ClickHouse テーブルに出力する過程で、コミュニティ版の ClickHouse はトランザクションをサポートしていないため、データが ClickHouse にシンクされる過程で Exactly-Once セマンティクスを保証できませんでした。この過程でジョブのフェールオーバーが発生すると、データが ClickHouse に重複して書き込まれます。

この課題に対処するため、BIGO ClickHouse は 2 フェーズコミットのトランザクション機構を実装しました。データを ClickHouse に書き込む際、まず書き込みモードを temporary に設定して、現在書き込んでいるデータが一時データであることを示します。挿入操作を実行して Insert ID を返し、その後 Insert ID に基づいて Commit 操作を実行すると、一時データが正式なデータに変換されます。

BIGO ClickHouse の 2 フェーズコミットトランザクション機構と Flink のチェックポイント機構を組み合わせて、ClickHouse Sink の Exactly-Once 書き込みセマンティクスを保証する ClickHouse Connector を実装しました。以下の通りです。

通常書き込み時、Connector は ClickHouse のシャードをランダムに選択して書き込みを行い、ユーザー設定に応じて単一コピーまたは二重コピーの insert 操作を実行し、書き込み後に Insert ID を記録します。この種の insert 操作が複数発生し、複数の Insert ID が生成されます。チェックポイント完了時に、これらの Insert ID を一括コミットして一時データを正式なデータに変換します。すなわち、2 つのチェックポイント間のデータ書き込みが完了します。
ジョブのフェールオーバー発生時、Flink ジョブのフェールオーバー再起動後、最新のチェックポイントからステートが復元されます。このとき、ClickHouse Sink の Operator State には前回タイムリーにコミットされなかった Insert ID が含まれている可能性があります。これらの Insert ID については、コミットをリトライします。データは既に ClickHouse に書き込まれているが、Insert ID が Operator State のデータに記録されていない場合、一時データであるため ClickHouse で照会されません。一定時間後、ClickHouse の期限切れクリーンアップ機構によりクリアされるため、ステートが最後のチェックポイントにロールバックされ、データが重複しないことが保証されます。
上記の機構により、Kafka から Flink 経由で計算後に ClickHouse に書き込まれる全プロセスのエンドツーエンド Exactly-Once セマンティクスが正常に保証され、データの重複や損失は発生しません。

2.2.3 プラットフォーム構築
BIGO ビッグデータプラットフォームのリアルタイムコンピューティングタスクをより適切に管理するため、同社は BIGO リアルタイムコンピューティングプラットフォーム Bigoflow を構築し、ユーザーに統一された Flink リアルタイムタスクアクセスを提供しています。プラットフォーム構築は以下の通りです。

Flink JAR、SQL、Python などの多彩なジョブタイプをサポートし、さまざまな Flink バージョンに対応して、社内のリアルタイムコンピューティング関連ビジネスの大半をカバーしています。
ワンストップ管理:ジョブの開発、投稿、実行、履歴表示、モニタリング、アラートを統合し、ジョブの実行状況をいつでも確認して問題を発見できます。
データリネージ:各ジョブのデータソース、データの用途、データ計算の流れを簡単に照会できます。

3. 適用シナリオ

3.1 OneSQL OLAP 分析プラットフォームの適用シナリオ
OneSQL OLAP 分析プラットフォームの社内での適用シナリオは、アドホッククエリへの適用です。以下の通りです。

ユーザーが Hue ページから投稿した SQL は、OneSQL バックエンドを通じて Flink SQL Gateway に転送され、Flink Session クラスターに送信されてクエリタスクが実行されます。Flink SQL Gateway はクエリタスクの実行進捗を取得して Hue ページに返し、クエリ結果を返します。

3.2 リアルタイムデータウェアハウスの適用シナリオ
リアルタイムデータウェアハウスの適用シナリオは、現在主に ABTest ビジネスです。以下の通りです。

ユーザーの元の行動ログデータは、Flink タスクで集約されてユーザー詳細データが生成され、ディメンションテーブルデータとフローディメンションテーブル JOIN を実行して、ClickHouse に出力されて多次元詳細ワイドテーブルが生成され、異なるディメンションに応じて集約されて異なるビジネスに適用されます。ABTest ビジネスを変革したことで、当該ビジネスの結果指標の生成時間を 8 時間前倒しし、リソース使用量を半分以上削減しました。

4. 今後の計画

OneSQL OLAP 分析プラットフォームと BIGO リアルタイムデータウェアハウスをより良く構築するため、リアルタイムコンピューティングプラットフォームの計画は以下の通りです。

Flink OLAP 分析プラットフォームを改善し、Hive SQL 構文サポートを拡充し、計算過程での JOIN データスキューの問題を解決します。

リアルタイムデータウェアハウスの構築を改善し、データレイク技術を導入して、リアルタイムデータウェアハウス内のタスクデータの再実行と小範囲でのトレースバックの問題を解決します。

Flink をベースにしたストリームバッチ統合データコンピューティングプラットフォームを構築します。

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.