すべてのプロダクト
Search
ドキュメントセンター

E-MapReduce:PySpark streaming on EMR Serverless Spark

最終更新日:Jul 08, 2026

ストリーム処理は、リアルタイムのビッグデータ分析に不可欠です。EMR Serverless Spark は、サーバー管理を不要にし、データ処理を簡素化する強力でスケーラブルなプラットフォームです。このトピックでは、EMR Serverless Spark を使用して PySpark ストリーミングジョブを送信する方法を説明し、ストリーム処理における使いやすさと保守性の高さを紹介します。

前提条件

ワークスペースが作成済みであること。詳細については、「ワークスペースの作成」をご参照ください。

操作手順

ステップ 1: Dataflow クラスターの作成とメッセージの生成

  1. EMR on ECS ページで、Kafka サービスを含むリアルタイム Dataflow クラスターを作成します。詳細については、「クラスターの作成」をご参照ください。

  2. EMR on ECS クラスターのマスターノードにログインします。詳細については、「クラスターへのログイン」をご参照ください。

  3. 次のコマンドを実行してディレクトリを変更します。

    cd /var/log/emr/taihao_exporter
  4. 次のコマンドを実行してトピックを作成します。

    # パーティション数 10、レプリケーション係数 2 で、taihaometrics という名前のトピックを作成します。
    kafka-topics.sh --partitions 10 --replication-factor 2 --bootstrap-server core-1-1:9092 --topic taihaometrics --create
  5. 次のコマンドを実行してメッセージを送信します。

    # kafka-console-producer を使用して、taihaometrics トピックにメッセージを送信します。
    tail -f metrics.log | kafka-console-producer.sh --broker-list core-1-1:9092 --topic taihaometrics

ステップ 2: ネットワーク接続の作成

  1. ネットワーク接続ページに移動します。

    1. EMR コンソールの左側メニューで、[EMR Serverless] > Spark を選択します。

    2. [Spark] ページで、対象のワークスペース名をクリックします。

    3. EMR Serverless Spark ページの左側メニューで、Normal Network Connection をクリックします。

  2. Normal Network Connection ページで、Create Network Connection をクリックします。

  3. Create Network Connection ダイアログボックスで、次のパラメータを設定し、OK をクリックします。

    パラメータ

    説明

    Name

    接続の名前を入力します。例: connection_to_emr_kafka。

    VPC

    EMR on ECS クラスターがデプロイされている VPC を選択します。

    使用可能な VPC がない場合は、Create VPC をクリックして VPC コンソールに移動し、VPC を作成します。詳細については、「VPC の作成と管理」をご参照ください。

    vSwitch

    EMR on ECS クラスターと同じ VPC 内の vSwitch を選択します。

    現在のゾーンに使用可能な vSwitch がない場合は、vSwitch をクリックして VPC コンソールに移動し、vSwitch を作成します。詳細については、「vSwitch の作成と管理」をご参照ください。

    StatusSucceeded になると、ネットワーク接続が作成されます。

ステップ 3: セキュリティグループルールの追加

  1. クラスターノードの vSwitch CIDR ブロックを取得します。

    Nodes ページで、ノードグループ名をクリックして関連する vSwitch を確認します。次に、VPC コンソールにログインし、vSwitch ページで vSwitch の CIDR ブロックを確認します。

    image

  2. セキュリティグループルールを追加します。

    1. Clusters ページで、対象のクラスターの ID をクリックします。

    2. 基本情報 ページで、Cluster Security Group の横にあるリンクをクリックします。

    3. [Security Group Details] ページの アクセスルール セクションで、[Add Rule] をクリックします。次のパラメータを設定し、OK をクリックします。

      パラメータ

      説明

      [Source]

      前のステップで取得した vSwitch CIDR ブロックを入力します。

      重要

      0.0.0.0/0 には設定しないでください。クラスターが外部アクセスにさらされます。

      [Destination (Current Instance)]

      ポート 9092 を入力します。

ステップ 4: JAR パッケージの OSS へのアップロード

kafka.zip を解凍し、アーカイブ内のすべての JAR パッケージを OSS にアップロードします。詳細については、「シンプルアップロード」をご参照ください。

ステップ 5: リソースファイルのアップロード

  1. EMR Serverless Spark ページの左側メニューで、Artifacts をクリックします。

  2. Artifacts ページで、Upload File をクリックします。

  3. Upload File ダイアログボックスで、アップロード領域をクリックし、pyspark_ss_demo.py ファイルを選択します。

ステップ 6: ストリーミングジョブの作成と開始

  1. EMR Serverless Spark ページの左側メニューで、Development をクリックします。

  2. Development タブで、image アイコンをクリックします。

  3. 名前を入力し、ジョブタイプに Application (Streaming) > PySpark を選択して、OK をクリックします。

  4. 新しい開発タブで、次のパラメータを設定し、その他の項目はデフォルト値のまま、Save をクリックします。

    パラメータ

    説明

    Main Python Resources

    前のステップの [Files] ページでアップロードした pyspark_ss_demo.py ファイルを選択します。

    Engine Version

    Spark バージョンを選択します。詳細については、「エンジンバージョン」をご参照ください。

    Execution Parameters

    クラスターの core-1-1 ノードの内部 IP アドレスを入力します。この IP は、Core ノードグループ内の Nodes ページで確認できます。

    Spark Configuration

    Spark 設定を指定します。以下は例です。

    spark.jars oss://path/to/commons-pool2-2.11.1.jar,oss://path/to/kafka-clients-2.8.1.jar,oss://path/to/spark-sql-kafka-0-10_2.12-3.3.1.jar,oss://path/to/spark-token-provider-kafka-0-10_2.12-3.3.1.jar
    spark.emr.serverless.network.service.name connection_to_emr_kafka
    説明
    • spark.jars :必要な外部 JAR パッケージの OSS パス。例のパスを、ステップ 4 でアップロードした JAR の OSS パスに置き換えてください。

    • spark.emr.serverless.network.service.name :ネットワーク接続の名前。例の値を、ステップ 2 で作成したネットワーク接続の名前に置き換えてください。

  5. Publish をクリックします。

  6. Publish ダイアログボックスで、OK をクリックします。

  7. ストリーミングジョブを開始します。

    1. Go to O&M をクリックします。

    2. START をクリックします。

ステップ 7: ログの表示

  1. Log Exploration タブをクリックします。

  2. Log Exploration タブで、アプリケーションの実行詳細と結果を確認します。

    image

関連トピック

PySpark 開発ワークフローの例については、「PySpark 開発クイックスタート」をご参照ください。