現在、私たちは毎秒ペタバイトのデータが生成される世界に住んでいます。そのため、データが生成され、さらに多くのデータが生成されるにつれて、ビジネス上の洞察をより正確に生成したいと考えている企業にとって、このデータをリアルタイムで分析および処理することは、ますます重要になります。
今日は、Spark Structured Streaming と Apache Kafka を使用して、架空の航空交通データに基づくリアルタイム データ分析を開発します。これらのテクノロジーが何であるかわからない場合は、この記事で取り上げる他の概念と同様に、それらのテクノロジーをより詳しく紹介するために私が書いた記事を読むことをお勧めします。ぜひチェックしてみてください。

Spark Structured Streaming と Apache Kafka を使用したリアルタイム データ処理の簡単な紹介
ゲハジ・アンク・2022年9月29日
私の GitHub で完全なプロジェクトをチェックアウトできます。
建築
データ エンジニアであるあなたが、航空交通に関するデータが毎秒生成される SkyX という航空会社で働いていると想像してください。
あなたは、海外で最も訪問される都市のランキングなど、これらのフライトからのリアルタイム データを表示するダッシュボードを開発するように依頼されました。ほとんどの人が離れる都市。そして世界中で最も多くの人を輸送する航空機です。
これは各フライトで生成されるデータです:
- aircraft_name: 航空機の名前。 SkyX では、利用可能な航空機は 5 機のみです。
- 出発地: 航空機が出発する都市。 SkyX は、世界中の 5 都市間のフライトのみを運航しています。
- 宛先: 航空機の目的地都市。前述したように、SkyX は世界中の 5 都市間のフライトのみを運航しています。
- 乗客: 航空機が輸送している乗客の数。すべての SkyX 航空機は、各フライトで 50 人から 100 人を乗せます。
以下は私たちのプロジェクトの基本的なアーキテクチャです:
- プロデューサー: 航空機の航空交通データを生成し、それを Apache Kafka トピックに送信する責任があります。
- コンシューマ: Apache Kafka トピックにリアルタイムで到着するデータのみを観察します。
- データ分析: Apache Kafka トピックに到着するデータをリアルタイムで処理および分析する 3 つのダッシュボード。最も多くの観光客を受け入れる都市の分析。ほとんどの人が他の都市を訪れるために出発する都市の分析。そして、世界中の都市間で最も多くの人を輸送する SkyX 航空機の分析。
開発環境の準備
このチュートリアルは、マシンに PySpark がすでにインストールされていることを前提としています。まだ確認していない場合は、ドキュメント内の手順を確認してください。
Apache Kafka については、Docker?? を介してコンテナ化して使用します。
そして最後に、仮想環境を通じて Python を使用します。
Docker 経由のコンテナ化による Apache Kafka
すぐに、skyx というフォルダーを作成し、その中にファイル docker-compose.yml を追加します。
$ mkdir skyx $ cd skyx $ touch docker-compose.yml
次に、次のコンテンツを docker-compose ファイル内に追加します。
version: '3.9' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 2181:2181 kafka: image: confluentinc/cp-kafka:latest depends_on: - zookeeper ports: - 29092:29092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
完了!これで、Kafka サーバーをアップロードできるようになりました。これを行うには、ターミナルに次のコマンドを入力します:
$ docker compose up -d $ docker compose ps
NAME COMMAND SERVICE STATUS PORTS skyx-kafka-1 "/etc/confluent/dock…" kafka running 9092/tcp, 0.0.0.0:29092->29092/tcp skyx-zookeeper-1 "/etc/confluent/dock…" zookeeper running 2888/tcp, 0.0.0.0:2181->2181/tcp, 3888/tcp
注: このチュートリアルでは、Docker Compose のバージョン 2.0 を使用しています。 docker と compose の間に「-」がないのはこのためです ☺.
ここで、プロデューサーによってリアルタイムで送信されたデータを保存するトピックを Kafka 内に作成する必要があります。これを行うには、コンテナ内の Kafka にアクセスしましょう:
$ docker compose exec Kafka bash
そして最後に、airtraffic というトピックを作成します。
$ kafka-topics --create --topic airtraffic --bootstrap-server localhost:29092
航空交通に関するトピックを作成しました。
仮想環境の作成
プロデューサー、つまりリアルタイムの航空交通データを Kafka トピックに送信する役割を担うアプリケーションを開発するには、kafka-python ライブラリを使用する必要があります。 kafka-python は、Apache Kafka と統合するプロデューサーとコンシューマーの開発を可能にするコミュニティ開発のライブラリです。
まず、requirements.txt というファイルを作成し、その中に次の依存関係を追加しましょう。
カフカパイソン
2 番目に、仮想環境を作成し、requirements.txt ファイルに依存関係をインストールします。
$ mkdir skyx $ cd skyx $ touch docker-compose.yml
完了!これで環境は開発の準備が整いました?
プロデューサーの育成
それでは、プロデューサーを作成しましょう。前述したように、プロデューサーは、新しく作成された Kafka トピックに航空交通データを送信する責任を負います。
アーキテクチャでも言われていましたが、SkyX は世界中の 5 つの都市間を飛行するだけで、利用可能な航空機は 5 台しかありません。各航空機には 50 人から 100 人が搭乗できることは注目に値します。
データはランダムに生成され、1 ~ 6 秒の間隔で JSON 形式でトピックに送信されることに注意してください。
行きましょう! src というサブディレクトリと、kafka という別のサブディレクトリを作成します。 kafka ディレクトリ内に、airtraffic_Producer.py というファイルを作成し、その中に次のコードを追加します。
version: '3.9' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 2181:2181 kafka: image: confluentinc/cp-kafka:latest depends_on: - zookeeper ports: - 29092:29092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
完了!私たちはプロデューサーを育成します。実行して、しばらく放置します。
$ python airtraffic_Producer.py
消費者開発
それでは、コンシューマを開発しましょう。これは非常に単純なアプリケーションになります。 Kafka トピックに到着したデータをターミナルにリアルタイムで表示するだけです。
引き続き kafka ディレクトリ内に、airtraffic_consumer.py というファイルを作成し、その中に次のコードを追加します。
$ docker compose up -d $ docker compose ps
ほら、とても簡単なことだと言いましたね。これを実行し、プロデューサーがトピックにデータを送信するときにリアルタイムで表示されるデータを確認します。
$ python airtraffic_consumer.py
データ分析: 観光客が最も多い都市
ここからはデータ分析を始めます。この時点で、最も観光客が多い都市のランキングをリアルタイムに表示するダッシュボード、つまりアプリケーションを開発します。つまり、データを to 列でグループ化し、passengers 列に基づいて合計を作成します。とてもシンプルです!
これを行うには、src ディレクトリ内に dashboards というサブディレクトリを作成し、tourists_analysis.py というファイルを作成します。次に、その中に次のコードを追加します:
$ mkdir skyx $ cd skyx $ touch docker-compose.yml
これで、spark-submit を通じてファイルを実行できるようになりました。でも、落ち着いてください! PySpark を Kafka と統合する場合は、spark-submit を別の方法で実行する必要があります。 --packages.
パラメータを通じて、Apache Kafka パッケージと Apache Spark の現在のバージョンを通知する必要があります。Apache Spark を Apache Kafka と統合するのが初めての場合、spark-submit の実行には時間がかかる場合があります。これは、必要なパッケージをダウンロードする必要があるためです。
データ分析をリアルタイムで確認できるように、プロデューサーがまだ実行されていることを確認してください。ダッシュボード ディレクトリ内で、次のコマンドを実行します:
$spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 Tourism_analysis.py
version: '3.9' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 2181:2181 kafka: image: confluentinc/cp-kafka:latest depends_on: - zookeeper ports: - 29092:29092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
データ分析: ほとんどの人が離れる都市
この分析は前の分析と非常に似ています。ただし、最も多くの観光客を受け入れる都市をリアルタイムで分析するのではなく、最も多くの人が離れる都市を分析します。これを行うには、leavers_analysis.py というファイルを作成し、その中に次のコードを追加します。
$ docker compose up -d $ docker compose ps
データ分析をリアルタイムで確認できるように、プロデューサーがまだ実行されていることを確認してください。ダッシュボード ディレクトリ内で、次のコマンドを実行します:
$spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 Leavers_analysis.py
NAME COMMAND SERVICE STATUS PORTS skyx-kafka-1 "/etc/confluent/dock…" kafka running 9092/tcp, 0.0.0.0:29092->29092/tcp skyx-zookeeper-1 "/etc/confluent/dock…" zookeeper running 2888/tcp, 0.0.0.0:2181->2181/tcp, 3888/tcp
データ分析: 最も多くの乗客を運ぶ航空機
この分析は前の分析よりもはるかに単純です。世界中の都市間で最も多くの乗客を輸送している航空機をリアルタイムで分析してみましょう。 aircrafts_analysis.py というファイルを作成し、その中に次のコードを追加します:
$ python -m venv venv $ venv\scripts\activate $ pip install -r requirements.txt
データ分析をリアルタイムで確認できるように、プロデューサーがまだ実行されていることを確認してください。ダッシュボード ディレクトリ内で、次のコマンドを実行します:
$spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0航空機_分析.py
$ mkdir skyx $ cd skyx $ touch docker-compose.yml
最終的な考慮事項
そして皆さん、ここで終わります!この記事では、Spark Structured Streaming と Apache Kafka を使用して、架空の航空交通データに基づくリアルタイム データ分析を開発します。
これを行うために、このデータをリアルタイムで Kafka トピックに送信するプロデューサーを開発し、その後、このデータをリアルタイムで分析するための 3 つのダッシュボードを開発しました。
気に入っていただければ幸いです。また次回お会いしましょう。
以上がSpark Structured Streaming と Apache Kafka を使用したリアルタイムの航空交通データ分析の詳細内容です。詳細については、PHP 中国語 Web サイトの他の関連記事を参照してください。

Pythonは、データサイエンス、Web開発、自動化タスクに適していますが、Cはシステムプログラミング、ゲーム開発、組み込みシステムに適しています。 Pythonは、そのシンプルさと強力なエコシステムで知られていますが、Cは高性能および基礎となる制御機能で知られています。

2時間以内にPythonの基本的なプログラミングの概念とスキルを学ぶことができます。 1.変数とデータ型、2。マスターコントロールフロー(条件付きステートメントとループ)、3。機能の定義と使用を理解する4。

Pythonは、Web開発、データサイエンス、機械学習、自動化、スクリプトの分野で広く使用されています。 1)Web開発では、DjangoおよびFlask Frameworksが開発プロセスを簡素化します。 2)データサイエンスと機械学習の分野では、Numpy、Pandas、Scikit-Learn、Tensorflowライブラリが強力なサポートを提供します。 3)自動化とスクリプトの観点から、Pythonは自動テストやシステム管理などのタスクに適しています。

2時間以内にPythonの基本を学ぶことができます。 1。変数とデータ型を学習します。2。ステートメントやループの場合などのマスター制御構造、3。関数の定義と使用を理解します。これらは、簡単なPythonプログラムの作成を開始するのに役立ちます。

10時間以内にコンピューター初心者プログラミングの基本を教える方法は?コンピューター初心者にプログラミングの知識を教えるのに10時間しかない場合、何を教えることを選びますか...

fiddlereveryversings for the-middleの測定値を使用するときに検出されないようにする方法

Python 3.6のピクルスファイルのロードレポートエラー:modulenotFounderror:nomodulenamed ...

風光明媚なスポットコメント分析におけるJieba Wordセグメンテーションの問題を解決する方法は?風光明媚なスポットコメントと分析を行っているとき、私たちはしばしばJieba Wordセグメンテーションツールを使用してテキストを処理します...


ホットAIツール

Undresser.AI Undress
リアルなヌード写真を作成する AI 搭載アプリ

AI Clothes Remover
写真から衣服を削除するオンライン AI ツール。

Undress AI Tool
脱衣画像を無料で

Clothoff.io
AI衣類リムーバー

AI Hentai Generator
AIヘンタイを無料で生成します。

人気の記事

ホットツール

AtomエディタMac版ダウンロード
最も人気のあるオープンソースエディター

ZendStudio 13.5.1 Mac
強力な PHP 統合開発環境

DVWA
Damn Vulnerable Web App (DVWA) は、非常に脆弱な PHP/MySQL Web アプリケーションです。その主な目的は、セキュリティ専門家が法的環境でスキルとツールをテストするのに役立ち、Web 開発者が Web アプリケーションを保護するプロセスをより深く理解できるようにし、教師/生徒が教室環境で Web アプリケーションを教え/学習できるようにすることです。安全。 DVWA の目標は、シンプルでわかりやすいインターフェイスを通じて、さまざまな難易度で最も一般的な Web 脆弱性のいくつかを実践することです。このソフトウェアは、

WebStorm Mac版
便利なJavaScript開発ツール

Safe Exam Browser
Safe Exam Browser は、オンライン試験を安全に受験するための安全なブラウザ環境です。このソフトウェアは、あらゆるコンピュータを安全なワークステーションに変えます。あらゆるユーティリティへのアクセスを制御し、学生が無許可のリソースを使用するのを防ぎます。
