Scenario solution demo based on real-time computing Flink version
リアルタイムアプリケーションログ分析
シナリオ説明
最初のシナリオの要件は比較的よくあるものです。このシナリオでは、車両プライバシー保護用の API を構築します。この API は、ユーザーがアップロードした車両写真に対してプライバシー保護処理を実行できるもので、ディープラーニングモデルに基づいています。
このモデルは API としてパッケージ化され、Alibaba Cloud のパブリッククラウド ECS にデプロイされて、世界中のユーザーがアクセスできます。この API でまず必要なのは、アクセス数やフィードバック頻度、ユーザーがどの国や地域からアクセスしているか、アクセスの特性はどのようなものか、攻撃なのか通常の利用なのかを分析することです。
このリアルタイム分析を実現するには、まず各サーバーに散在する各 API の大量のリアルタイムアプリケーションログを収集できる必要があります。収集だけでなく、比較的タイムリーかつリアルタイムに処理できる必要もあります。処理には、ディメンションテーブルのクエリやウィンドウ集計などが含まれます。これらはストリーミングコンピューティングでは比較的一般的な操作です。最後に、これらの処理結果を高スループットで低レイテンシな環境に配置することで、下流の分析システムがリアルタイムでデータにアクセスできるようになります。
この一連のフローは複雑ではありませんが、非常に重要な機能を表しています。すなわち、Flink を代表とするリアルタイムコンピューティングを活用することで、ビジネス意思決定者に秒単位でデータドリブンな意思決定機能を提供できるということです。
デモソリューションアーキテクチャ
このデモの実装方法を見ていきましょう。このアーキテクチャにはいくつかの重要なポイントがあります。
まず、右上角は構築された API 環境で、Flask と Python に主流の NGINX と Gunicorn を組み合わせて API として構築しています。この API をコンテナイメージ化し、イメージを通じて Alibaba Cloud の ECS にデプロイする必要があります。高同時実行性と低レイテンシのため、第 7 層ロードバランサーも設置され、さらに前面には API Gateway が配置されて、ユーザーが API を呼び出す機能を提供しています。
同時に、このデモでは Web アプリも提供しており、ユーザーはコードだけでなくグラフィカルなインターフェースからも API を利用できます。フロントエンドのユーザーが API を呼び出すと、SLS の Simple Log Service を使用して API サーバーからリアルタイムにアプリケーションログを収集し、簡単な処理を施してからリアルタイムコンピューティング Flink に配信します。
Flink の優れた点は、Simple Log Service からのログ配信をサブスクライブし、ストリーミングコンピューティングの形でこのログに対してウィンドウ集計やディメンションテーブルのクエリ結合などの操作を実行できることであり、さらに SQL を使って複雑なビジネスロジックをカスタマイズできるという利点もあります。
すべてのデータ処理が完了した後、Flink はストリームデータを Hologres にテーブル形式の構造化データとして書き込みます。Hologres は単なるデータストアではなく、下流の BI データを支える OLAP 系のエンジンでもあります。これらを組み合わせることで、ビッグデータのリアルタイムログ収集と分析のフレームワークが完成します。
ソリューション分析
各部分の使用方法を詳しく見ていきましょう。
車両プライバシー API をリアルタイム分析のデータソースとして使用する
Web アプリを通じて、ユーザーは車両写真を簡単にアップロードでき、API がぼかし処理を施します。スクリーンレコーディングでは、API による処理後に写真の背景がぼやけ、ナンバープレート部分やその他のプライバシー情報も隠されているのが確認できます。
SLS Log Center
ユーザーがこの API にアクセスすると、バックグラウンドの Simple Log Service がリアルタイムでログを収集します。
ログ収集後、Logtail のデータ変換・処理機能を活用して、元のログに対して一定の解析と変換を行います。これには IP アドレスの国や都市などの地理情報への解析(緯度精度を含む)が含まれ、下流の分析処理でこの情報を活用できるようになります。また、基本的なサービスに加え、強力なグラフィカルデータ分析機能も提供しています。
リアルタイムコンピューティング Flink 版
ここでは一次データ分析やデータサーベイ機能を実行でき、元のログの変換が下流のビジネス支援の要件を満たしているかを確認できます。ログの収集、変換、処理が完了した後、ログは Log Hub 配信を通じてストリーム処理センター、つまりリアルタイムコンピューティング Flink に転送されます。
実際には、「配信」という用語は正確ではありません。実際には、Flink が Log Hub 内の Log Store に保存された処理済みログ情報を積極的にサブスクライブしています。Flink の非常に優れた点は、SQL を使ってビジネスロジックを記述できることです。条件分岐の変換や処理なども含め、SQL を記述して「オンライン化」をクリックするだけで、Flink ジョブとしてパッケージ化され、Flink クラスター内にホスティングされます。クラスター内では、このコンソールから非常に便利にアクセスできます。
では、現在のクラスターの使用状況はどうでしょうか。CPU の状態はどうか、異常はないか、エラー報告はないか、配信状況の確認など、すべて Flink 経由で直接ホスティング・モニタリングできます。これは非常に大きな利点で、運用管理の作業をほぼ気にする必要がありません。
Hologres (HSAP)
Flink の処理完了後、ストリームデータは Flink が提供するインターフェイスを通じて、ストレージシステム Hologres にテーブル形式の構造化データとして直接書き込まれます。Hologres の特に大きな特徴は、OLTP と OLAP の両方を兼ね備えていることです。
具体的には、OLTP として高速なデータ書き込みができると同時に、書き込まれたデータに対して高同時実行かつ低レイテンシのクエリ分析を実行できます。つまり、OLAP エンジンの機能を併せ持っており、この 2 つを統合したものであるため、Hologres は HSAP とも呼ばれています。
image34.png
DataV ダッシュボード
このアーキテクチャでは、主に処理済みデータを下流のエンドユーザーやビジネス意思決定者に表示するために使用されます。意思決定者は消費用のリアルタイム大画面を確認できます。
このリアルタイム大画面は、API へのアクセスに伴い、秒単位の遅延で最新の処理済み情報を反映します。この DataV のリアルタイム大画面により、意思決定者がデータを目にするまでの遅延を大幅に削減できます。
従来のバッチ処理方式を使用した場合、処理ごとにテラバイト単位のデータが必要で、処理時間は数時間に及ぶ可能性があります。Flink を中心としたエンドツーエンドのリアルタイムコンピューティングソリューションを採用すれば、この遅延を数時間から数秒、場合によっては 1 秒以内に圧縮できます。
車両エンジンのリアルタイム予知保全
シナリオ説明
2 番目のビジネスシナリオは、IoT のシミュレートされたテレメトリデータを組み合わせて、道路上を走行する車両のエンジンに異常の兆候がないかを分析・判定し、事前に問題の可能性を判断するものです。放置すれば、3 か月後に何らかの部品が故障する可能性があるといった予測ができます。これも実運用シナリオで頻繁に求められる要件で、予知保全と呼ばれます。実際のアプリケーションシナリオでは、予知保全はお客様のコストを大幅に削減できます。なぜなら、問題が発生してから修理するよりも、故障前に事前に交換する方がはるかに効果的だからです。
デモソリューションアーキテクチャ
このような現実に比較的近いシナリオを実現するため、車載装置の診断システム OBD II に着目しました。これには典型的なデータが含まれており、その一部を収集・処理してシミュレーションを行います。より現実に近い車両エンジンの動作データをシミュレートするプログラムを作成しました。
今回は実際に車両を道路上で走行させることができないため、このシミュレーションプログラムを使用して、さまざまな統計分析手法により走行データをシミュレートし、できるだけ現実に近い効果を実現しています。
このプログラムはシミュレートされた走行エンジンテレメトリデータを Kafka に配信し、リアルタイムコンピューティング Flink で Kafka トピックを消費・サブスクライブして、各トピックに応じた異なるストリーム処理を実行します。処理結果の一部は OSS にアーカイブされ、保存することで履歴データとして取得できます。もう一方は、開発済みの異常検知モデルにリアルタイムデータソースとして直接配信され、PAI-EAS にデプロイされて Flink から直接呼び出せます。
その後、機械学習による判定を経て、現在のエンジンデータに異常の兆候があるかを確認し、その結果をデータベースに書き込んで下流のシステムが順次消費できるようにします。Flink でリアルタイム処理されたデータの一部は OSS にアーカイブされます。
このアーカイブデータは、実際には履歴データとしてモデルの構築や再学習に利用されます。一定期間ごとに走行特性に変化(データドリフトとして知られる)が生じた場合、新しく生成された履歴データを使用してモデルを再学習でき、再学習済みモデルを Web サービスとして PAI-EAS にデプロイし、Flink から呼び出せるようにします。こうして Lambda ベースのビッグデータソリューションが完成します。
ソリューション分析
シミュレーション走行データの生成
まず、シミュレーションデータを生成する作業が必要です。エンジンのテレメトリデータを模擬した OBD データをシミュレートし、クラウドに配信して分析に利用します。ここでは Function Compute を使用します。Function Compute は非常に便利なフルマネージドサービスです。
次に、ローカルで開発した Python スクリプトを直接 Function Compute にコピーして設定でき、このマネージドコンピューティング環境でシミュレーションデータ生成スクリプトを実行できるので、非常に便利です。
このデモでは、Function Compute が 1 分ごとに実行され、1 バッチのテレメトリデータを生成し、毎回 3 秒間隔で Kafka に配信して、実際の環境のデータをできるだけ忠実に再現した頻度でデータ生成を行います。
走行データの収集・配信
Kafka はビッグデータで一般的に使用される Pub/Sub システムで、柔軟性と優れたスケーラビリティを備えています。Alibaba Cloud 上の Kafka では、EMR に Kafka クラスターを構築することも、Message Queue for Apache Kafka というフルマネージドサービスを利用して完全なサービスを構築することもできます。
これは Kafka システムです。このデモでは便宜上、Kafka でフルマネージドの Pub/Sub システムを構築しています。このシステムは、前述の車両から配信されたエンジンテレメトリデータを保存するためにのみ使用されます。実際の本番環境では、1 台の車両だけでなく、数万、数十万台の車両を扱う必要があります。Kafka を使用すれば、非常に便利にスケーリングできます。フロントエンドの車両数が 10 台でも 10 万台でも、全体的なアーキテクチャを大きく変更する必要がなく、これらのスケーラビリティの要件に柔軟に対応できます。
リアルタイム計算と異常分析モデルの呼び出し
リアルタイムコンピューティング部分では、引き続き Flink のリアルタイムコンピューティングシステムを使用しますが、このデモでは Blink 専用クラスターを使用しています。これはセミマネージド型のリアルタイムコンピューティングプラットフォームです。実際には、前のシナリオのフルマネージド方式とほぼ同じです。
デモ作成当時、一部の地域で Flink のフルマネージド版がまだ提供されていなかったため、Blink 専用クラスターというサービスを選択しました。これもリアルタイムコンピューティングファミリーに属しており、使い方は Flink フルマネージド版とほぼ同じです。開発者はビジネスロジック処理のスクリプトを書くことに集中し、「オンライン化」をクリックするだけで、残りは基本的に Flink に完全に管理されます。異常がないかのモニタリングやチューニングなどの作業も非常に便利です。
ここで特筆すべきは、PAI-EAS のモデル呼び出しインターフェイスが Flink に組み込まれていることです。Flink がストリームデータをリアルタイム処理しながら、一部データを PAI に送信してモデルの推論を実行でき、その結果をリアルタイムストリームデータと組み合わせて、最終的に下流のストレージシステムに書き込みます。これは Flink コンピューティングプラットフォームの拡張性とスケーラビリティを体現しています。
異常検知モデルの開発
この部分では、グラフィカルな学習プラットフォームを使用して、非常にシンプルな二値分類モデルを設計・開発する方法を示しています。
この二値分類モデルは主に、過去のエンジン履歴データから、どのような特徴がエンジンの問題を示し、どのような値が正常範囲であるかを学習するものです。このモデルにより、新たに生成されるエンジンデータに対して判定を行う基準ができ、ビジネス担当者が現在のエンジンのデータ問題を事前に予測するのに役立ちます。
モデルのデプロイと呼び出しサービス
このモデルは過去のデータから関連する特徴とデータパターンを学習しています。モデル開発の全プロセスで使用する studio は、ドラッグアンドドロップで完全に構築でき、コードをほとんど書くことなく、ボタン操作だけでモデル開発を実現できます。非常に便利で迅速です。さらに素晴らしいのは、モデル開発完了後、PAI のワンクリックデプロイで REST API および Web サービスとしてパッケージ化し、PAI プラットフォーム上に配置して他のユーザーが呼び出せるようにできることです。ワンクリックデプロイ後、デプロイされたモデルのサービスに対してテスト呼び出しを行うのも非常に便利です。
高スループット構造化データストア (RDS)
モデルのデプロイが完了したら、Flink に異常の有無を判定させます。ストリームデータはリアルタイム処理された後、最終的に MySQL データベースに書き込まれます。
このデータベースはデータソースとして、下流のリアルタイム大画面にデータ支援を提供します。これにより、ビジネス担当者はリアルタイムで、つまり数秒ごとに、現在道路上を走行中の車両に問題がないかを確認できます。
ニアリアルタイムダッシュボード
DataV の大画面は 5 秒ごとに更新されるように設定されています。つまり、5 秒ごとにデータベースから最新のテレメトリデータと異常の有無を判定するデータが取得され、大画面に表示されます。
赤色はその時点で収集されたデータに問題があることを示し、青色は正常、つまり比較的正常なデータを表します。このデータの正常性基準は、前述のシミュレーションデータを生成する Function Compute のロジックによって完全に制御されています。Function Compute のロジック内で、エンジンに異常があるように見せるデータが人為的に追加されており、このデモの異常部分がより明確に反映されるようになっています。
シナリオ説明
最初のシナリオの要件は比較的よくあるものです。このシナリオでは、車両プライバシー保護用の API を構築します。この API は、ユーザーがアップロードした車両写真に対してプライバシー保護処理を実行できるもので、ディープラーニングモデルに基づいています。
このモデルは API としてパッケージ化され、Alibaba Cloud のパブリッククラウド ECS にデプロイされて、世界中のユーザーがアクセスできます。この API でまず必要なのは、アクセス数やフィードバック頻度、ユーザーがどの国や地域からアクセスしているか、アクセスの特性はどのようなものか、攻撃なのか通常の利用なのかを分析することです。
このリアルタイム分析を実現するには、まず各サーバーに散在する各 API の大量のリアルタイムアプリケーションログを収集できる必要があります。収集だけでなく、比較的タイムリーかつリアルタイムに処理できる必要もあります。処理には、ディメンションテーブルのクエリやウィンドウ集計などが含まれます。これらはストリーミングコンピューティングでは比較的一般的な操作です。最後に、これらの処理結果を高スループットで低レイテンシな環境に配置することで、下流の分析システムがリアルタイムでデータにアクセスできるようになります。
この一連のフローは複雑ではありませんが、非常に重要な機能を表しています。すなわち、Flink を代表とするリアルタイムコンピューティングを活用することで、ビジネス意思決定者に秒単位でデータドリブンな意思決定機能を提供できるということです。
デモソリューションアーキテクチャ
このデモの実装方法を見ていきましょう。このアーキテクチャにはいくつかの重要なポイントがあります。
まず、右上角は構築された API 環境で、Flask と Python に主流の NGINX と Gunicorn を組み合わせて API として構築しています。この API をコンテナイメージ化し、イメージを通じて Alibaba Cloud の ECS にデプロイする必要があります。高同時実行性と低レイテンシのため、第 7 層ロードバランサーも設置され、さらに前面には API Gateway が配置されて、ユーザーが API を呼び出す機能を提供しています。
同時に、このデモでは Web アプリも提供しており、ユーザーはコードだけでなくグラフィカルなインターフェースからも API を利用できます。フロントエンドのユーザーが API を呼び出すと、SLS の Simple Log Service を使用して API サーバーからリアルタイムにアプリケーションログを収集し、簡単な処理を施してからリアルタイムコンピューティング Flink に配信します。
Flink の優れた点は、Simple Log Service からのログ配信をサブスクライブし、ストリーミングコンピューティングの形でこのログに対してウィンドウ集計やディメンションテーブルのクエリ結合などの操作を実行できることであり、さらに SQL を使って複雑なビジネスロジックをカスタマイズできるという利点もあります。
すべてのデータ処理が完了した後、Flink はストリームデータを Hologres にテーブル形式の構造化データとして書き込みます。Hologres は単なるデータストアではなく、下流の BI データを支える OLAP 系のエンジンでもあります。これらを組み合わせることで、ビッグデータのリアルタイムログ収集と分析のフレームワークが完成します。
ソリューション分析
各部分の使用方法を詳しく見ていきましょう。
車両プライバシー API をリアルタイム分析のデータソースとして使用する
Web アプリを通じて、ユーザーは車両写真を簡単にアップロードでき、API がぼかし処理を施します。スクリーンレコーディングでは、API による処理後に写真の背景がぼやけ、ナンバープレート部分やその他のプライバシー情報も隠されているのが確認できます。
SLS Log Center
ユーザーがこの API にアクセスすると、バックグラウンドの Simple Log Service がリアルタイムでログを収集します。
ログ収集後、Logtail のデータ変換・処理機能を活用して、元のログに対して一定の解析と変換を行います。これには IP アドレスの国や都市などの地理情報への解析(緯度精度を含む)が含まれ、下流の分析処理でこの情報を活用できるようになります。また、基本的なサービスに加え、強力なグラフィカルデータ分析機能も提供しています。
リアルタイムコンピューティング Flink 版
ここでは一次データ分析やデータサーベイ機能を実行でき、元のログの変換が下流のビジネス支援の要件を満たしているかを確認できます。ログの収集、変換、処理が完了した後、ログは Log Hub 配信を通じてストリーム処理センター、つまりリアルタイムコンピューティング Flink に転送されます。
実際には、「配信」という用語は正確ではありません。実際には、Flink が Log Hub 内の Log Store に保存された処理済みログ情報を積極的にサブスクライブしています。Flink の非常に優れた点は、SQL を使ってビジネスロジックを記述できることです。条件分岐の変換や処理なども含め、SQL を記述して「オンライン化」をクリックするだけで、Flink ジョブとしてパッケージ化され、Flink クラスター内にホスティングされます。クラスター内では、このコンソールから非常に便利にアクセスできます。
では、現在のクラスターの使用状況はどうでしょうか。CPU の状態はどうか、異常はないか、エラー報告はないか、配信状況の確認など、すべて Flink 経由で直接ホスティング・モニタリングできます。これは非常に大きな利点で、運用管理の作業をほぼ気にする必要がありません。
Hologres (HSAP)
Flink の処理完了後、ストリームデータは Flink が提供するインターフェイスを通じて、ストレージシステム Hologres にテーブル形式の構造化データとして直接書き込まれます。Hologres の特に大きな特徴は、OLTP と OLAP の両方を兼ね備えていることです。
具体的には、OLTP として高速なデータ書き込みができると同時に、書き込まれたデータに対して高同時実行かつ低レイテンシのクエリ分析を実行できます。つまり、OLAP エンジンの機能を併せ持っており、この 2 つを統合したものであるため、Hologres は HSAP とも呼ばれています。
image34.png
DataV ダッシュボード
このアーキテクチャでは、主に処理済みデータを下流のエンドユーザーやビジネス意思決定者に表示するために使用されます。意思決定者は消費用のリアルタイム大画面を確認できます。
このリアルタイム大画面は、API へのアクセスに伴い、秒単位の遅延で最新の処理済み情報を反映します。この DataV のリアルタイム大画面により、意思決定者がデータを目にするまでの遅延を大幅に削減できます。
従来のバッチ処理方式を使用した場合、処理ごとにテラバイト単位のデータが必要で、処理時間は数時間に及ぶ可能性があります。Flink を中心としたエンドツーエンドのリアルタイムコンピューティングソリューションを採用すれば、この遅延を数時間から数秒、場合によっては 1 秒以内に圧縮できます。
車両エンジンのリアルタイム予知保全
シナリオ説明
2 番目のビジネスシナリオは、IoT のシミュレートされたテレメトリデータを組み合わせて、道路上を走行する車両のエンジンに異常の兆候がないかを分析・判定し、事前に問題の可能性を判断するものです。放置すれば、3 か月後に何らかの部品が故障する可能性があるといった予測ができます。これも実運用シナリオで頻繁に求められる要件で、予知保全と呼ばれます。実際のアプリケーションシナリオでは、予知保全はお客様のコストを大幅に削減できます。なぜなら、問題が発生してから修理するよりも、故障前に事前に交換する方がはるかに効果的だからです。
デモソリューションアーキテクチャ
このような現実に比較的近いシナリオを実現するため、車載装置の診断システム OBD II に着目しました。これには典型的なデータが含まれており、その一部を収集・処理してシミュレーションを行います。より現実に近い車両エンジンの動作データをシミュレートするプログラムを作成しました。
今回は実際に車両を道路上で走行させることができないため、このシミュレーションプログラムを使用して、さまざまな統計分析手法により走行データをシミュレートし、できるだけ現実に近い効果を実現しています。
このプログラムはシミュレートされた走行エンジンテレメトリデータを Kafka に配信し、リアルタイムコンピューティング Flink で Kafka トピックを消費・サブスクライブして、各トピックに応じた異なるストリーム処理を実行します。処理結果の一部は OSS にアーカイブされ、保存することで履歴データとして取得できます。もう一方は、開発済みの異常検知モデルにリアルタイムデータソースとして直接配信され、PAI-EAS にデプロイされて Flink から直接呼び出せます。
その後、機械学習による判定を経て、現在のエンジンデータに異常の兆候があるかを確認し、その結果をデータベースに書き込んで下流のシステムが順次消費できるようにします。Flink でリアルタイム処理されたデータの一部は OSS にアーカイブされます。
このアーカイブデータは、実際には履歴データとしてモデルの構築や再学習に利用されます。一定期間ごとに走行特性に変化(データドリフトとして知られる)が生じた場合、新しく生成された履歴データを使用してモデルを再学習でき、再学習済みモデルを Web サービスとして PAI-EAS にデプロイし、Flink から呼び出せるようにします。こうして Lambda ベースのビッグデータソリューションが完成します。
ソリューション分析
シミュレーション走行データの生成
まず、シミュレーションデータを生成する作業が必要です。エンジンのテレメトリデータを模擬した OBD データをシミュレートし、クラウドに配信して分析に利用します。ここでは Function Compute を使用します。Function Compute は非常に便利なフルマネージドサービスです。
次に、ローカルで開発した Python スクリプトを直接 Function Compute にコピーして設定でき、このマネージドコンピューティング環境でシミュレーションデータ生成スクリプトを実行できるので、非常に便利です。
このデモでは、Function Compute が 1 分ごとに実行され、1 バッチのテレメトリデータを生成し、毎回 3 秒間隔で Kafka に配信して、実際の環境のデータをできるだけ忠実に再現した頻度でデータ生成を行います。
走行データの収集・配信
Kafka はビッグデータで一般的に使用される Pub/Sub システムで、柔軟性と優れたスケーラビリティを備えています。Alibaba Cloud 上の Kafka では、EMR に Kafka クラスターを構築することも、Message Queue for Apache Kafka というフルマネージドサービスを利用して完全なサービスを構築することもできます。
これは Kafka システムです。このデモでは便宜上、Kafka でフルマネージドの Pub/Sub システムを構築しています。このシステムは、前述の車両から配信されたエンジンテレメトリデータを保存するためにのみ使用されます。実際の本番環境では、1 台の車両だけでなく、数万、数十万台の車両を扱う必要があります。Kafka を使用すれば、非常に便利にスケーリングできます。フロントエンドの車両数が 10 台でも 10 万台でも、全体的なアーキテクチャを大きく変更する必要がなく、これらのスケーラビリティの要件に柔軟に対応できます。
リアルタイム計算と異常分析モデルの呼び出し
リアルタイムコンピューティング部分では、引き続き Flink のリアルタイムコンピューティングシステムを使用しますが、このデモでは Blink 専用クラスターを使用しています。これはセミマネージド型のリアルタイムコンピューティングプラットフォームです。実際には、前のシナリオのフルマネージド方式とほぼ同じです。
デモ作成当時、一部の地域で Flink のフルマネージド版がまだ提供されていなかったため、Blink 専用クラスターというサービスを選択しました。これもリアルタイムコンピューティングファミリーに属しており、使い方は Flink フルマネージド版とほぼ同じです。開発者はビジネスロジック処理のスクリプトを書くことに集中し、「オンライン化」をクリックするだけで、残りは基本的に Flink に完全に管理されます。異常がないかのモニタリングやチューニングなどの作業も非常に便利です。
ここで特筆すべきは、PAI-EAS のモデル呼び出しインターフェイスが Flink に組み込まれていることです。Flink がストリームデータをリアルタイム処理しながら、一部データを PAI に送信してモデルの推論を実行でき、その結果をリアルタイムストリームデータと組み合わせて、最終的に下流のストレージシステムに書き込みます。これは Flink コンピューティングプラットフォームの拡張性とスケーラビリティを体現しています。
異常検知モデルの開発
この部分では、グラフィカルな学習プラットフォームを使用して、非常にシンプルな二値分類モデルを設計・開発する方法を示しています。
この二値分類モデルは主に、過去のエンジン履歴データから、どのような特徴がエンジンの問題を示し、どのような値が正常範囲であるかを学習するものです。このモデルにより、新たに生成されるエンジンデータに対して判定を行う基準ができ、ビジネス担当者が現在のエンジンのデータ問題を事前に予測するのに役立ちます。
モデルのデプロイと呼び出しサービス
このモデルは過去のデータから関連する特徴とデータパターンを学習しています。モデル開発の全プロセスで使用する studio は、ドラッグアンドドロップで完全に構築でき、コードをほとんど書くことなく、ボタン操作だけでモデル開発を実現できます。非常に便利で迅速です。さらに素晴らしいのは、モデル開発完了後、PAI のワンクリックデプロイで REST API および Web サービスとしてパッケージ化し、PAI プラットフォーム上に配置して他のユーザーが呼び出せるようにできることです。ワンクリックデプロイ後、デプロイされたモデルのサービスに対してテスト呼び出しを行うのも非常に便利です。
高スループット構造化データストア (RDS)
モデルのデプロイが完了したら、Flink に異常の有無を判定させます。ストリームデータはリアルタイム処理された後、最終的に MySQL データベースに書き込まれます。
このデータベースはデータソースとして、下流のリアルタイム大画面にデータ支援を提供します。これにより、ビジネス担当者はリアルタイムで、つまり数秒ごとに、現在道路上を走行中の車両に問題がないかを確認できます。
ニアリアルタイムダッシュボード
DataV の大画面は 5 秒ごとに更新されるように設定されています。つまり、5 秒ごとにデータベースから最新のテレメトリデータと異常の有無を判定するデータが取得され、大画面に表示されます。
赤色はその時点で収集されたデータに問題があることを示し、青色は正常、つまり比較的正常なデータを表します。このデータの正常性基準は、前述のシミュレーションデータを生成する Function Compute のロジックによって完全に制御されています。Function Compute のロジック内で、エンジンに異常があるように見せるデータが人為的に追加されており、このデモの異常部分がより明確に反映されるようになっています。
Related Articles
-
A detailed explanation of Hadoop core architecture HDFS
Knowledge Base Team
-
What Does IOT Mean
Knowledge Base Team
-
6 Optional Technologies for Data Storage
Knowledge Base Team
-
What Is Blockchain Technology
Knowledge Base Team
Explore More Special Offers
-
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
