DeltaLake's practical sharing in industrial brain

2020 年の Industrial Brain 3.0 のリリース以来、Industrial Brain は長年の発展を遂げてきました。本記事では、Industrial Brain の構築における DeltaLake の優れた実践について共有します。主な内容は以下の通りです。

(1) 異なる場所からの異種ストリームメッセージの処理

(2) ストリームバッチ融合のデータ分析

(3) トランザクション処理とアルゴリズムサポート

1. 異なる場所からの異種ストリームメッセージの処理

工業企業にとって、データソースは世界中に分散していることが一般的です。グループレベルのユーザーは、各地のデータセンターからデータを集約して取得したいと考えています。

DeltaLake と構造化ストリーミングを組み合わせることで、以下の 2 つの処理を実現できます。

1. 各工場エリアの Kafka リアルタイムデータを集約し、データプラットフォームの Kafka に書き込んでリアルタイムデータアプリケーションに対応する

2. 各工場エリアの Kafka データをデータプラットフォームの HDFS にアーカイブし、オフライン分析型データアプリケーションに対応する

上記のタスクを完了できるビッグデータコンポーネントは Flink や Flume など多数ありますが、最終的に DeltaLake を選択した理由は以下の通りです。

1. 複数の Kafka トピックの正規表現による消費をサポート

SubscribePattern を使用すると、正規表現ルールで複数のトピックから同時にデータを消費でき、工業団地内に消費対象のトピックが多い場合に非常に便利です。

2. HDFS のサポートとスモールファイル統合のカプセル化

「Kafka のデータをリアルタイムで HDFS に書き込む」シナリオでは、DeltaLake を使用すると非常に便利です。主な理由は 2 つあります。

1. HDFS への書き込みをネイティブでサポートするため、Flink で HDFS シンカーを作成する手間や、Flume クラスターの追加運用保守の手間を解消できます。

2. ストリーミングウェアハウジングの各シナリオは、データアーキテクトにとってパフォーマンスと適時性のトレードオフです。Flink、Flume、Spark Streaming のいずれを使用する場合でも、「ローリング書き込み容量 (または件数) しきい値」と「ローリング書き込み時間しきい値」が設計されています。実際の導入プロセスでは、データレイテンシとパフォーマンスに関するビジネス要件の違いに応じて、この 2 つが衡量されます。たとえば、レイテンシ許容度が低いシナリオでは、容量または件数のしきい値を非常に小さく (場合によっては 1 に) 設定して、新しいデータを素早く書き込むことができます。ただし、その副作用としてシンカーの IO が頻繁になり、HDFS に多数のスモールファイルが生成され、データの読み書き性能や DataNode のパフォーマンスに影響を与えます。レイテンシ許容度が高いシナリオでは、配信エンジニアは件数しきい値と時間しきい値を大きく設定して IO 性能を向上させることが多く、データレイテンシを犠牲にします。これは一般的な方法ですが、実際の生産プロセスでは、多数のストリームジョブに対して多数の異なる設定を維持するコストは無視できません。

DeltaLake を使用すると、この処理がはるかに容易になります。すべてのストリームジョブのローリング書き込みしきい値を同じ値 (たとえば、すべて小さく) に設定することで、すべてのストリームジョブがより良いデータレイテンシを得られます。同時に、DeltaLake の機能である Optimize と Vacuum を組み合わせて、定期タスクを設定して周期的に実行し、スモールファイルをマージまたは削除して HDFS の性能を確保することで、データ開発全体の作業をよりシンプルにし、運用保守を改善できます。

Optimize 機能の参照

2. ストリームバッチ融合のデータ分析

生産製造プロセスにおいて、機械装置の安定稼働は完成品の品質にとって極めて重要です。装置が安定しているかどうかを判断する最も直感的な方法は、特定のセンサーの長期履歴トレンドを確認することです。実際のプロジェクト導入時、配信エンジニアはストリーム処理を使用して、Kafka 内の大量のセンサー時系列データを OLAP ストレージ (Alibaba Cloud ADB、TSDB、HBase など) に処理して書き込み、上位のデータ分析アプリケーションの高同時実行性と低応答時間のリアルタイムクエリ要件をサポートします。

ただし、実際の状況はこれよりもはるかに複雑です。工業企業の情報化とデジタル化のレベルは一般的に高くなく、業界ごとに生産プロセスの自動化レベルもまちまちであるため、装置のリアルタイムデータの多くは実際には正確ではなく、人手による介入や再計算を経た後に使用できるものが多数あります。

そのため、実際の導入プロセスでは、「ローリング上書き」モードを使用して OLAP ストレージ内のデータを持続的に書き換え、OLAP を「リアルタイム増分エリア」と「定期カバレッジエリア」に分割することがよく行われます。以下の図に示す通りです。

上記の図は 1 つの OLAP ストレージで、すべてのデータはオレンジ色と青色の部分に分けられています。上位のデータアプリケーションはこれら 2 つのエリアのデータを区別なくクエリできます。唯一の違いは、最新のオレンジ色データはストリームコンピューティングジョブによって Kafka からリアルタイムで取得され、処理後に書き込まれるのに対し、青色エリアは履歴データが定期的に計算 (修正ロジックを追加) されて書き込まれ、昨日以前のリアルタイムデータを修正する点です。このようにしてサイクルを繰り返し、データの適時性を確保すると同時に、履歴データを修正・カバレッジしてデータの正確性を保証します。

従来は、ストリーム + バッチの Lambda アーキテクチャを使用し、2 つの異なるコンピューティングエンジンを使用してストリームとバッチを処理することが多く行われていました。以下の図に示す通りです。

Lambda アーキテクチャの欠点もここから見て取れます。異なるプラットフォーム上で 2 つのコードを保守し、その計算ロジックが完全に一致していることを確認するのは手間のかかる作業です。DeltaLake の導入後、事態は比較的シンプルになります。Spark のストリーミングとバッチの統合設計は、コードの再利用とクロスプラットフォームのロジック統一の問題をうまく解決しています。DeltaLake の特性 (ACID、OPTIMIZE など) と組み合わせることで、この作業をより洗練された形で実行できます。以下の図に示す通りです。

また、ストリーミングとバッチの統合は Spark 独自の機能ではありませんが、Alibaba Cloud E-MapReduce (EMR) は SparkSQL と Spark Streaming の上に SQL をカプセル化しており、ビジネス担当者が Flink SQL に似た構文を使用して低しきい値でジョブ開発を行えるようにしています。これにより、ストリーミングとバッチのシナリオでのコード再利用と運用保守作業が容易になり、プロジェクト配信効率の向上において大きな意義があります。詳細については EMR ドキュメントを参照してください。

3. トランザクション処理とアルゴリズムサポート

従来のデータウェアハウスでは、モデル構築プロセスでトランザクションを導入することはほとんどありませんでした。データウェアハウスはデータの変化を反映する必要があるため、ACID を使用してデータウェアハウスとビジネスシステムの整合性を保つのではなく、遅変ディメンションなどの方法を使用してデータ状態の変化を記録することが多く行われていました。

ただし、工業データプラットフォームの導入においては、トランザクションには独自のユースケースがあります。たとえば、生産スケジューリングは、すべての工業企業が注目する重要な課題です。生産スケジューリングはグループレベルで実施されることが多く、顧客の注文、材料在庫、工場の能力に基づいて当期の生産需要を合理的に分解・配置し、能力の合理的な配分を実現します。一方、スケジューリングはより微視的です。工場レベルでは、作業指示書、材料、実際の生産状況に基づいて生産計画をリアルタイムで動的に調整し、リソースの利用率を最大化します。どちらも多くのデータ融合を必要とする計画問題です。以下の図に示す通りです。

生産スケジューリングアルゴリズムに必要な生データは、注文と計画データを提供する ERP、材料データを提供する WMS、作業指示書と工程データを提供する MES など、複数の業務システムから取得されることが多いです。これらのデータは (物理的および論理的に) 統合されて初めて、生産スケジューリングアルゴリズムの有効な入力として使用できます。したがって、導入プロセスでは、各システムのデータを格納するための統合ストレージが必要です。同時に、生産スケジューリングアルゴリズムはデータの有効性に対しても一定の要件があり、すべての業務システムと可能な限り整合したデータを入力として、その時点の生産状況を正確に反映し、より良いスケジューリングを実現する必要があります。

従来、このシナリオは以下のように処理していました。

1) 各業務システムの CDC 機能を使用するか、個別のプログラムを作成してデータの変更を準リアルタイムでポーリングして取得する

2) リレーショナルデータベースに書き込み、このプロセスでデータマージのロジックを処理し、リレーショナルデータベース内のデータと業務システムデータを準リアルタイムで整合させる

3) 生産スケジューリングエンジンがトリガーされると、RDB からデータをプルして計算を実行する

このアーキテクチャにはいくつかの明らかな課題があり、主なものは以下の通りです。

1) RDB をビッグデータストレージの代わりに使用し、計算時にデータをメモリにクエリしてロードする方式は、大規模なデータ量に対して非常に困難です

2) Hive エンジンを使用して中間 RDB を置き換える場合、Hive 3.x では ACID がサポートされていますが、リアルタイム性能と MapReduce プログラミングフレームワークのアルゴリズム (ソルバー) サポートがエンジニアリング要件を満たせません

現在、DeltaLake を導入し、Spark の機能と組み合わせてアーキテクチャを最適化しようとしています。

最適化されたアーキテクチャの利点は以下の通りです。

1) HDFS + Spark を使用して RDB を中間ストレージに置き換え、データ量が大規模な場合のストレージ課題を解決

2) Spark Streaming + DeltaLake を使用して元データを連携し、DeltaLake の ACID 特性を使用してデータが中間ストレージに入る際のマージロジックを処理し、ストリーミング時にデータのマージ + 最適化を同時に行い、読み書き性能を確保

3) 生産スケジューリングエンジンは、中間プラットフォームからメモリにクエリデータを転送して計算するのではなく、アルゴリズムタスクを Spark ジョブとしてカプセル化してコンピューティングプラットフォームに送信して計算を完了します。これにより、Spark ML プログラミングフレームワークのアルゴリズムと Python の優れたサポート、および Spark の分散コンピューティング機能を活用して、複数回の反復が必要な計画アルゴリズムを分散処理できます

4) DeltaLake のタイムトラベル機能を使用してデータバージョンを管理またはロールバックでき、アルゴリズムモデルのデバッグと評価に非常に有益です

4. まとめ

1) DeltaLake のコア機能である ACID は、データのリアルタイム性と正確性への要件が高いアプリケーション、特にアルゴリズムアプリケーションにとって非常に有益で、Spark の ML に対するネイティブサポートをより効果的に活用できます

2) DeltaLake の Optimize + Vacuum とストリーミングのデータウェアハウジング機能を組み合わせることで、上流の Kafka データを大量に連携する際に優れた互換性を持ち、運用保守コストを効果的に削減できます

3) Alibaba Cloud EMR チームがカプセル化した Streaming SQL 開発ストリームジョブは、大規模データセンタープロジェクト導入時の開発しきい値とコストを効果的に削減できます

現在、Industrial Brain における DeltaLake の適用はまだ実験段階にあります。ストリーミングウェアハウジング、生産スケジューリングエンジン、ストリームバッチ融合などの複数のシナリオが、Industrial Brain の複数のプロジェクトで適用されています。同時に、これらのシナリオは徐々に Industrial Brain の標準プロダクトとなっており、Industrial Brain 3.0 のデータ + アルゴリズムシーンの視覚的編集機能と複製機能を組み合わせることで、ディスクリート製造、自動車、鉄鋼などの業界のシナリオに迅速に複製し、AI 機能を活用して中国の産業に利益をもたらすことができます。

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.