Construction and application of feature platform in Shuhe

1、特徴量プラットフォームの概要

まず、特徴量プラットフォームの概要を説明します。
特徴量プラットフォーム全体は 4 つのレイヤー、すなわちデータサービス、ストレージサービス、計算エンジン、原始ストレージに分かれています。
データサービスレイヤーは外部サービスを提供し、主に以下の 4 種類を含みます。


・1 つ目は、従来の API ポイントクエリです。

・2 つ目は、ポーリングです。

・3 つ目は、イベントメッセージです。

・4 つ目は、同期呼び出し計算です。


同期呼び出し計算サービスはリアルタイムコンピューティングであり、オンサイトのポリシー計算に相当します。
一方、API チェックサービスは事前計算して保存する方式です。
データサービスを提供するために、特徴量行ストレージと特徴量列ストレージの 2 つのサービスモードが用意されており、それぞれ API ポイントクエリと範囲クエリをサポートします。
計算エンジンはオフラインコンピューティングエンジンとストリームバッチコンピューティングエンジンの 2 種類があります。
特徴量プラットフォームの最下層はオフラインコンピューティングを支える原始ストレージであり、イベントストレージはストリームバッチコンピューティングをサポートします。


次に、MySQL を例に特徴量プラットフォームの簡略化されたデータフローを説明します。


まずオフライン部分では、MySQL データベースのデータを Sqoop などの抽出ツールで EMR に抽出し、Hive 処理を経て最終的な演算結果を HBase と ClickHouse に保存します。
これらはそれぞれ特徴量行ストレージと特徴量列ストレージに対応し、API ポイントクエリと範囲クエリサービスを提供します。
同時に、MySQL の Binlog が Kafka にリアルタイムで書き込まれ、Kafka のデータは Flink ストリームバッチ統合計算エンジンで消費されます。
また、Kafka のデータはイベントストレージである HBase にも消費され、イベントストレージ HBase のデータも Flink ストリームバッチ統合計算エンジンに提供されます。
エンジンでの計算後、データは HBase と ClickHouse に書き込まれ、イベントメッセージも送信されます。
イベントストア HBase に転送されたデータは、リアルタイム呼び出しサービスを提供できます。


2、特徴量ストレージサービス

次に、特徴量ストレージサービスについて説明します。


特徴量は以下の 4 つのカテゴリに分類されます。


・同期特徴量:リアルタイム書き込み、オフライン補正、ストリームバッチ処理を行います。

・即時計算特徴量:API 呼び出しによる演算で、オフラインバッチ計算とロジックの一貫性を保ちます。

・リアルタイム特徴量:従来のリアルタイムリンクで複雑なリアルタイムロジックを実現でき、ストリームバッチに置き換え可能です。

・オフライン特徴量:従来のオフラインリンクで複雑なオフラインロジックを実現します。


オフラインリンクが不可欠な理由は以下のとおりです。


1 つ目として、リアルタイムリンクは純粋な増分リンクです。
リンクが非常に長くなり、どの环节でも問題が発生する可能性があります。
エラーが発生すると、長期間にわたりデータが自動補正されません。


2 つ目として、リアルタイムリンクには適時性が求められます。
特にマルチストリーム結合の場合、遅延が発生するとできる限り早くフォールバック結果を返す必要があります。
リアルタイム特徴量の最終的なエラー率を制御し、エラーを小規模な期間に限定するために、オフラインリンクによる補正が必要です。
特徴量ストレージサービスの変更方法は 2 つあります。
1 つは同期特徴量で、自身の変更を持つストリームバッチリンクです。
その他の特徴量は、一般的にリアルタイム + オフライン + リアルタイムコンピューティングの組み合わせで処理されます。


MySQL を例に、ストレージサービスの全体的なデータフローを説明します。
オフライン部分では、MySQL データベースのデータを Sqoop 抽出ツールで EMR に抽出し、Hive 処理を経て最終的な演算結果を HBase と ClickHouse に保存します。
同時に、Binlog が Kafka にリアルタイムで書き込まれ、Kafka のデータは Flink ストリームバッチコンピューティングエンジンで消費されます。
エンジンの計算後、データは HBase と ClickHouse に書き込まれます。
HBase と ClickHouse は API ポイントクエリと範囲クエリサービスを提供します。


2.1 リアルタイム特徴量のデータフロー

リアルタイムデータストリームでは、MySQL が Binlog を通じて Kafka にデータを書き込み、その他のイベントクラスの Kafka データも同様に処理されます。
計算後、結果を別の Kafka に書き込み、最終的に消費データを HBase と ClickHouse に書き込みます。


2.2 オフライン特徴量のデータフロー

オフライン特徴量のデータフローでは、MySQL は Sqoop を通じて抽出され、OSS は Spark またはその他の方法で抽出され、Kafka は Flume を使用して抽出し EMR に取り込まれます。
その後、Hive または Spark で処理を実行し、同時に HBase と ClickHouse に書き込みます。


2.3 同期特徴量のデータフロー

同期特徴量のデータフローでは、MySQL の Binlog がリアルタイム Kafka に書き込まれ、Kafka のデータがリアルタイムでイベントストレージに書き込まれます。
同時に、MySQL のオフライン補正と初期化も実行されます。
Flink はストリームとバッチを同時に処理し、HBase と ClickHouse に書き込みます。


2.4 リアルタイム計算特徴量のデータフロー

リアルタイムコンピューティング特徴量のデータフローでは、HBase と ClickHouse のデータに基づいて API ポイントクエリと範囲クエリサービスを提供します。


以上がストレージサービス全体の説明です。
この部分には特徴量ストレージサービスの大部分が含まれており、図のオレンジ色の部分に示されています。


3、ストリームバッチ統合ソリューション

特徴量ストレージサービスのみを提供していた時期に、多くの課題やビジネス上の要望が見つかりました。
まず、いくつかの質問を紹介します。


・既存のモデル戦略を集中的に活用しても、まだ使用されていないデータはないか?
たとえば、MySQL 内の状態変化のタイムスタンプデータなどです。


・入力特徴量のオフラインロジックは十分に完全か?
なぜリアルタイム入力特徴量を再設計して補完する必要があるのか?
オフライン入力特徴量をリアルタイム化するにはロジックを再構築する必要があります。
その中には、従来の方法ではリアルタイム変換が複雑すぎて完了できないものもあります。


・使用シナリオが不確定で、ポイントクエリとバッチ実行を区別できない場合、両方を同時にカバーできないか?
多くのビジネス担当者にとって、必要なモデルや戦略が最終的にバッチ実行とポイントクエリのどちらを必要とするかは未知です。
両方の要件を満たす方法はないでしょうか。


・ストリーム処理ロジックは理解しにくい。
なぜストリーム結合が必要なのか?
直接「フェッチ」できないのか?
モデル開発者にとってストリーム処理のフローは不明瞭であり、リアルタイム特徴量の作成が困難です。


・リアルタイムモデル戦略のバックテストに時間がかかるが、短縮できないか?


・モデル戦略の開発トレーニングは速いが、リアルタイム入力特徴量のオンライン開発には時間がかかる。
高速化できないか?


これらの課題に対し、以下のソリューションを提案します。


・[データ] 状態変化データを保存し、任意の時点のデータスライス状態を復元できます。
これには追加の利点もあります。
ストリームバッチ統合スキームでモデルトレーニングを行う際、未来のデータを取得できないため、特徴量トラバーサル問題が発生しません。


・[ロジック] ストリームとバッチを統合し、ストリームに注力して一貫したロジックを実現します。
キャリバーの検証は不要です。
このデータはトレーニングに使用され、同じデータがオンラインとバックテストにも使用されるため、最終結果の一貫性を保証できます。


・[実行] ストリーム、バッチ、呼び出しを統合し、さまざまなシナリオに柔軟に対応します。


・[開発] ストリームマージの代わりに「フェッチ」を使用し、リアルタイムストリーム固有の概念をカプセル化してリアルタイム開発のハードルを下げます。


・[テスト] 任意の期間のバックトラッキングテストをサポートし、リアルタイム開発とテストの速度を向上します。


・[リリース] セルフサービスのストリームバッチ統合モデルを開発して直接オンラインにリリースし、コミュニケーションリンクを削減してオンライン効率を向上します。


従来のリアルタイムストリーミングスキームには Lambda と Kappa があります。


Lambda はリアルタイムとオフラインの 2 組のロジックを提供し、最終的にデータベースで結合します。
Lambda の利点は、アーキテクチャがシンプルで、オフラインバッチ処理とリアルタイムストリーム処理の組み合わせが良く、リアルタイムコンピューティングコストが安定して制御しやすく、オフラインデータの補正が容易な点です。
欠点は、リアルタイムとオフラインのデータ結果の整合性維持が困難で、2 つのシステムをメンテナンスする必要がある点です。


3.1 ストリームバッチ統合スキーム

Kappa はリアルタイムロジックで既存データを保存し、毎回 1 つのデータスライスを取得して最終的にマージします。
Kappa の利点は、リアルタイム処理モジュールのみをメンテナンスすればよく、オフラインとリアルタイムのデータ統合が不要でメッセージをリプレイできる点です。
欠点は、メッセージミドルウェアのキャッシング機能に強く依存し、リアルタイムデータ処理中にデータ損失が発生する可能性があり、金融分野では許容できない点です。


Kappa はオフライン処理モジュールを廃止したため、オフラインコンピューティングがより信頼性が高く安定しているという特徴も放棄しています。
Lambda はオフラインコンピューティングの安定性を保証しますが、デュアルシステムのメンテナンスコストは非常に高く、2 組のコードの運用保守は非常に複雑です。


そこで、Lambda + Kappa のストリームバッチ統合スキームを提案します。
図に示すように、データフローの前半は Lambda アーキテクチャで、中心は HBase イベントストレージです。
後半はユーザーがストリーム処理とバッチ処理を完了する Kappa アーキテクチャです。


上の図は MySQL を例にした全体的なストリームバッチソリューションを示しています。
まず、MySQL の Binlog が Kafka に入り、オフライン補正とスライスを通じてイベントセンターにデータを送信します。
同時に、同じ Kafka を使用してリアルタイムストリームをトリガーします。
次に、イベントセンターはデータフェッチとオフラインバッチ実行サービスを提供します。
最後に、メタデータセンターがデータを一元的に管理・メンテナンスし、同期の問題を回避します。
Flink が全体のロジックサービスを提供します。


3.2 イベントセンター

図のイベントセンターは Lambda アーキテクチャを使用してすべての変更データを保存し、日次補正を実行します。
コールドホット混合ストレージとリヒーティングメカニズムを通じて最適なコストパフォーマンスを追求します。
また、Flink のウォーターマークメカニズムを参照して現在値の同期を確保します。
さらに、イベントセンターはメッセージ転送メカニズムと非同期から同期への変換メカニズムを提供し、ストリーム結合を「フェッチ」に置き換えます。
トリガーによるメッセージ受信とトリガーによるポーリング呼び出しをサポートし、インターフェイスにバックトラッキング機能を提供します。


以下にイベントセンターのデータフローを説明します。
図に示すように、MySQL や Kafka などの複数のデータソースが異なるパスを通じて Kafka に転送されます。
次に、Flink が Kafka を直接消費し、HBase ホットストレージにリアルタイムで書き込みます。
さらに、オフライン補正されたデータも EMR を通じて HBase ホットストレージに書き込まれます。
別のレプリカメカニズムが HBase ホットストレージと HBase コールドストレージ間のレプリケーションを完了します。
HBase コールドストレージのデータも HBase ホットストレージにリヒーティングされます。


イベントセンター全体のストレージ構造は図に示すとおりです。
コールドストレージには主体データのみが格納されます。
ホットストレージには主体データに加えて、異なるインデックス用の 3 つのテーブルがあります。
ホットストレージの TTL は通常 32 日で、特殊な状況下では調整可能です。


イベントセンターの読み取りデータストリームでは、リアルタイムトリガーは Kafka から、バックトラッキングとフェッチはいずれも HBase ホットストレージから行われます。
内部のリヒーティングメカニズムは HBase コールドストレージから HBase ホットストレージにデータを更新します。
この部分のロジックは開発者に対して透明であり、データの来源を意識する必要はありません。


以下にイベントセンターのウォーターマークメカニズムとストリーム結合について説明します。


2 つのストリームを結合するとします。
これは外部キーで関連付けられた 2 つのテーブルと単純に理解できます。
いずれかのテーブルが変更されると、少なくとも 1 回の最終的な完全な結合レコードをトリガーする必要があります。
2 つのフローをそれぞれ A と B とし、A が先に到着すると想定します。
イベントセンターのウォーターマークメカニズムが有効な場合、ストリーム A の現在イベントはストリーム A がトリガーされた時点で既にイベントセンターに記録されています。
2 つのケースがあります。


ストリーム B の関連データがイベントセンターから取得できる場合、ストリーム A の現在イベントがイベントセンターに記録された時点からストリーム B の現在イベントが実行されて関連データが読み取られる時点までの間に、ストリーム B がイベントセンターへの記録を完了しており、この時点のデータは完全です。


ストリーム B の関連データがイベントセンターから取得できない場合、イベントセンターのウォーターマークメカニズムはストリーム B の関連イベントがまだトリガーされていないことを示します。
ストリーム A の現在イベントは既にイベントセンターに書き込まれているため、ストリーム B の関連イベントがトリガーされる際にストリーム A の現在イベントデータは必ず取得可能であり、この時点のデータは完全です。
したがって、イベントセンターのウォーターマークメカニズムにより、ストリーム結合を「フェッチ」に置き換えた後、少なくとも 1 回の計算で完全なデータが確保されることが保証されます。


トリガーによるメッセージ受信はメッセージ転送を通じて完了します。
外部システムがリクエストを送信すると、Kafka に転送され、Kafka のデータが同時にイベントセンターに入ります。
次に、対応する計算がトリガーされます。
最後に、メッセージキューを使用して計算結果が送信され、外部システムがこのメッセージの結果を受信します。


同時に、ポーリングサービスとメッセージ転送も提供します。
前半のメカニズムはメッセージ受信メカニズムと同じですが、計算結果を再保存してサービスを提供するイベントセンターが追加されています。


3.3 アクセスの一貫性

データ取得のもう 1 つの課題は、メタベースから直接データを取得しない限り最新のデータを得ることができない点です。
ただし、この操作はメインデータベースに負荷をかけるため、一般的に禁止されています。
データの一貫性を確保するために、以下の対策を講じています。
まず、一貫性を 4 つのレベルに分類します。


結果整合性:更新されたデータは一定時間後にアクセス可能になります。
ストリームバッチ統合スキーム全体はデフォルトで結果整合性を保証します。


トリガーフローの強力な整合性 (遅延可):トリガーフローのフェッチ処理において、現在データより前のデータが取得可能であることを保証します。
ウォーターマークスキームを使用し、ウォーターマークが満たされない場合は遅延します。


データ取得の強力な整合性 (遅延可):ユーザーの時間要件より前のデータが取得可能であることを保証します。
ウォーターマークスキームを使用し、ウォーターマークが満たされない場合は遅延します。


データ取得の強力な整合性 (遅延なし):ユーザーの時間要件より前のデータが取得可能であることを保証します。
ウォーターマークが満たされない場合、データソースから増分的に直接補完します。
増分データ取得はデータソースに負荷をかけます。


3.4 ストリームバッチ統合運用

ストリームバッチ統合運用には PyFlink を使用します。
Python を使用するのは、モデルやポリシー開発者が Python に慣れているためであり、Flink がロジックの一貫性を保証します。
PyFlink に基づいて、複雑なトリガーロジック、複雑なアクセスロジックをカプセル化し、コードフラグメントを再利用できます。


PyFlink のコード構成は図に示すとおりで、トリガー、メインロジック、出力の 3 つの部分を含みます。
これら 3 つの部分を自分で実装する必要はなく、カプセル化された出力を選択するだけです。


Flink の全体的なデータフローもシンプルです。
最上部にトリガーロジックがあり、メインロジックがトリガーされます。
メインロジック内にはデータフェッチロジックがあり、フェッチを完了してから出力ロジックを実行します。
ここで、トリガーロジック、フェッチロジック、出力ロジックの基盤となるカプセル化は、ストリームバッチの変化に柔軟に対応するため、入出力を不変に保ちます。
ほとんどの場合、ロジック自体はストリームバッチ環境の変化を考慮する必要がありません。


以下に PyFlink の代表的な使用フローを紹介します。
まずトリガーフローを選択し、フェッチと前処理ロジックを記述します。
公開済みのフェッチまたは処理ロジックコードをインポートし、サンプリングロジックとテスト実行を設定してテスト実行結果を取得します。
その後、分析プラットフォームでさらなる分析とトレーニングを行います。
トレーニング完了後、ジョブ内でトレーニング済みモデルを選択し、必要に応じて初期化関連パラメータを設定します。
最後に、モデルをオンラインにリリースします。


4、モデルポリシー呼び出しスキーム

以下の 4 つの呼び出しスキームを提供します。


特徴量ストレージサービススキーム:Flink ジョブが事前計算を実行し、計算結果を特徴量ストレージサービスプラットフォームに書き込み、データサービスプラットフォームを通じて外部サービスを提供します。


インターフェーストリガー - ポーリングスキーム:イベントセンターのメッセージ転送インターフェイスを呼び出してポーリングし、Flink ジョブが計算結果を返すまで待機します。


インターフェーストリガー - メッセージ受信スキーム:イベントセンターのメッセージ転送インターフェイスを呼び出して Flink ジョブの実行をトリガーし、Flink ジョブが返す実行結果メッセージを受信します。


直接メッセージ受信スキーム:Flink ジョブが返す実行結果メッセージを直接受信します。


4.1 特徴量ストレージサービススキーム

特徴量ストレージサービスは、リアルタイム、オフライン補正、オフライン初期化の 3 つのケースに分かれます。
新しい変数がオンラインになる、または既存の変数のロジックが変更される場合、全量データを 1 回更新する必要があり、この場合はオフライン初期化が必要です。
リアルタイムフローはリアルタイムでトリガーされ、オフライン補正とオフライン初期化はバッチでトリガーされます。
フェッチロジックがある場合、データは HBase からフェッチされます。
リアルタイムとオフラインのジョブではフェッチプロセスが異なりますが、カプセル化されているため開発者が意識する必要はありません。
リアルタイム Flink ジョブの結果は Kafka に送信され、オフライン補正とオフライン初期化の結果は EMR に取り込まれ、最終的に特徴量ストレージ、すなわち HBase と ClickHouse に書き込まれます。


上の図は特徴量ストレージサービススキームのタイミングを示しています。
Kafka が Flink をトリガーし、Flink の演算結果が特徴量ストレージに書き込まれます。
トリガー時に外部呼び出しがあっても最新データは取得できません。
演算が完了してストレージに書き込まれるまで、更新されたデータは取得できません。


4.2 インターフェーストリガー - ポーリングスキーム

インターフェーストリガーポーリングスキームでは、トリガー呼び出しがメッセージ転送をトリガーし、Kafka に転送されます。
そして Flink が演算結果を Kafka に出力します。
この時間が単一リクエストの時間を超えない場合、直接返します。
この場合、トリガーポーリングは 1 回の呼び出しに縮退します。
逆に、時間を超える場合はイベントストア HBase にデータを書き込み、ポーリング呼び出しで結果を取得します。


インターフェーストリガー - ポーリングスキームのタイミングチャートは上の図に示すとおりです。
外部呼び出しがトリガーされると、メッセージ転送が Flink の計算をトリガーします。
Flink の計算とデータベースへの書き込みプロセス中に複数回のポーリングが行われます。
固定時間内に取得できない場合はタイムアウトを通知します。
次のポーリングでデータが既に書き込まれていれば、取得成功です。


4.3 インターフェーストリガー - メッセージ受信スキーム

インターフェーストリガー - メッセージ受信スキームはポーリングの簡略版です。業務システムがメッセージ受信をサポートしている場合、リンク全体がよりシンプルになります。メッセージ転送サービスを通じて計算をトリガーし、結果メッセージをリッスンするだけです。

インターフェーストリガー - メッセージ受信スキームのタイミングはシリアルです。トリガー後に Flink の計算が実行され、計算完了後にメッセージ受信メカニズムを通じて結果データが呼び出し側に送信されます。

4.4 直接メッセージ受信スキーム

直接メッセージ受信スキームは純粋なストリーミングプロセスです。Kafka が Flink の計算をトリガーし、計算後にデータがメッセージキューに転送され、相手がサブスクライブして受信します。全体のタイミングも非常にシンプルで、下の図に示すとおりです。

データの利用方法は、即時呼び出し、リアルタイムストリーム、オフラインバッチデータの 3 つのカテゴリに分類されます。適時性は順に低下します。これら 3 つのケースをイベントセンターを通じて一緒に登録します。最後に、イベントセンターを中継として Flink に関連データを提供するだけで済みます。Flink 側では、呼び出し方法を意識する必要がありません。

最後に、ストリーム、バッチ、呼び出しの統合に関する 4 つのスキームをまとめます。

・特徴量ストレージサービススキーム:特徴量ストレージサービスを通じて永続的な特徴量ストレージを提供します。API ポイントクエリと特徴量ポーリングサービスを提供します。

・インターフェーストリガー - ポーリングスキーム:イベントセンターのメッセージ転送とメッセージクエリサービスを通じて同期呼び出し計算サービスを提供します。

・インターフェーストリガー - メッセージ受信スキーム:イベントセンターのメッセージ転送サービスを通じてイベントメッセージサービスを提供します。

・直接メッセージ受信スキーム:複雑なイベントトリガーをサポートし、イベントメッセージを提供します。

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.