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

E-MapReduce:Kafka へのデータの読み書き

最終更新日:Nov 09, 2025

このトピックでは、EMR Serverless Spark でストリーミングタスクを開発および実行して、Kafka へのデータの読み書きを行う方法について説明します。Java アーカイブ (JAR) パッケージをアップロードし、ストリーミングタスクを作成して実行し、ログまたはコンソールを確認して結果を検証します。

前提条件

  • Serverless Spark ワークスペースが作成されていること。 詳細については、「ワークスペースの作成」をご参照ください。

  • Message Queue for Kafka インスタンスが作成されていること。

    このトピックでは、クラウドベースの Message Queue for Kafka インスタンスを例として使用します。 詳細については、「インスタンスの購入とデプロイ」をご参照ください。

手順

ステップ 1: Kafka 関連の JAR を OSS にアップロードする

Serverless Spark エンジンには、バージョン esr-2.8.0、esr-3.4.0、esr-4.4.0 以降の組み込み Kafka JAR パッケージが含まれています。このステップはスキップできます。 別のエンジンバージョンを使用する場合は、次のステップを実行する必要があります。

  1. Spark のバージョンに基づいて、次の依存関係を手動で追加できます。

  2. ファイルを解凍した後、すべての JAR パッケージを OSS にアップロードします。このトピックでは、kafka-spark35-jars.zip を例として使用します。

    hadoop fs -put /root/spark-sql-kafka-0-10_2.12-3.5.3.jar oss://<YOUR_BUCKET>.<region>.oss-dls.aliyuncs.com/
    hadoop fs -put /root/kafka-clients-3.4.1.jar oss://<YOUR_BUCKET>.<region>.oss-dls.aliyuncs.com/
    hadoop fs -put /root/spark-token-provider-kafka-0-10_2.12-3.5.3.jar oss://<YOUR_BUCKET>.<region>.oss-dls.aliyuncs.com/
    hadoop fs -put /root/commons-pool2-2.11.1.jar oss://<YOUR_BUCKET>.<region>.oss-dls.aliyuncs.com/
    OSS コンソールまたは別の方法を使用してファイルをアップロードすることもできます。

ステップ 2: ネットワーク接続を作成する

Serverless Spark は、Kafka サービスにアクセスするために Kafka ネットワークに接続できる必要があります。ネットワーク接続の詳細については、「EMR Serverless Spark と他の VPC 間のネットワーク接続」をご参照ください。

ステップ 3: テストファイルを準備する

次のセクションでは、Scala コードの例を示します。その他の例については、「Kafka 統合に関する Spark ドキュメント」をご参照ください。SparkExample-1.0-SNAPSHOT-jar-with-dependencies.jar をダウンロードしてテスト JAR パッケージを直接使用するか、自分でコードをパッケージ化することができます。pom.xml ファイルの構成については、「付録」をご参照ください。

コード例

Kafka からの読み取り

このコードは、Message Queue for Kafka の Topic にサブスクライブし、メッセージを [key, value] のペアとしてコンソールに出力します。

object StreamingReadKafka {

  def main(args: Array[String]): Unit = {

    import org.apache.spark.sql.SparkSession

    val servers = args(0)
    val topic = args(1)

    val spark = SparkSession.builder()
      .appName("test read kafka")
      .getOrCreate()

    import spark.implicits._

    val df = spark
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", servers)
      .option("subscribe", topic)
      .load()

    val query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .as[(String, String)]
      .writeStream
      .format("console")
      .start()

    query.awaitTermination()
  }
}

Kafka への書き込み

このコードは、組み込みの Spark データソースからデータを読み取り、値とタイムスタンプを含むレコードを生成し、value 列を key 列と value 列の両方として使用し、毎秒 1 レコードのレートで Kafka にレコードを書き込みます。

object StreamingWriteKafka {

  def main(args: Array[String]): Unit = {

    import org.apache.spark.sql.SparkSession
    import org.apache.spark.sql.functions._
    import org.apache.spark.sql.streaming.Trigger

    val servers = args(0)
    val topic = args(1)
    val checkpointDir = args(2)

    val spark = SparkSession.builder()
      .appName("test write kafka")
      .getOrCreate()

    val df = spark.readStream
      .format("rate")
      .option("rowsPerSecond", "1")
      .load()
      .withColumn("key", col("value"))
      .withColumn("value", col("value"))

    val query = df
      .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
      .writeStream
      .format("kafka")
      .option("kafka.bootstrap.servers", servers)
      .option("topic", topic)
      .option("checkpointLocation", checkpointDir)
      .trigger(Trigger.ProcessingTime("10 seconds"))
      .start()

    query.awaitTermination()
  }
}

ステップ 4: JAR パッケージをアップロードする

  1. ファイルアップロードページに移動します。

    1. E-MapReduce コンソールにログインします。

    2. 左側のナビゲーションウィンドウで、[EMR Serverless] > [Spark] を選択します。

    3. [Spark] ページで、ターゲットワークスペースをクリックします。

    4. EMR Serverless Spark ページで、左側のナビゲーションウィンドウの [ファイル管理] をクリックします。

  2. [ファイル管理] ページで、[ファイルのアップロード] をクリックします。

  3. [ファイルのアップロード] ダイアログボックスで、アップロードエリアをクリックするか、パッケージをエリアにドラッグして JAR パッケージを選択できます。

ステップ 5: ストリーミングタスクを作成して実行する

  1. EMR Serverless Spark ページで、左側のナビゲーションウィンドウの [データ開発] をクリックします。

  2. [開発] タブで、新規アイコン image をクリックします。

  3. 表示されるダイアログボックスで、[名前] を入力し、[ストリーミングタスク] の JAR タイプを選択して、[OK] をクリックします。

  4. 右上隅で、ターゲットキューを選択します。

    キューの追加方法の詳細については、「リソースキューの管理」をご参照ください。

  5. 新しい開発タブで、次のパラメーターを構成し、他のパラメーターはデフォルト設定のままにします。次に、[公開] をクリックします。

    パラメーター

    説明

    メイン JAR リソース

    前のステップでアップロードした JAR パッケージを選択します。この例では、パッケージは SparkExample-1.0-SNAPSHOT-jar-with-dependencies.jar です。

    エンジンバージョン

    適切な Spark バージョンを選択します。この例では、バージョンは esr-4.3.0 です。

    メインクラス

    Spark タスクを送信するときに指定するメインクラス。

    • Kafka からの読み取りの例: org.example.StreamingReadKafka

    • Kafka への書き込みの例: org.example.StreamingWriteKafka

    引数

    メインクラスに渡されるカスタムパラメーター。複数のパラメーターはスペースで区切ります。

    • Kafka エンドポイント情報。例: alikafka-serverless-cn-xxxxxx-1000-vpc.alikafka.aliyuncs.com:9092

    • サブスクライブする Kafka Topic。例: test

    • チェックポイントを保存するパス: oss://<YOUR_BUCKET_PATH>/。このパラメーターは、Kafka への書き込みにのみ必要です。

    ネットワーク接続

    ステップ 2 で作成したネットワークを選択します。

    Spark 構成

    spark.emr.serverless.user.defined.jars パラメーターを使用して、Kafka 関連の JAR パッケージを指定します。複数の JAR はカンマで区切ります。

    spark.emr.serverless.user.defined.jars  oss://<YOUR_BUCKET_PATH>/commons-pool2-2.11.1.jar,oss://<YOUR_BUCKET_PATH>/kafka-clients-3.4.1.jar,oss://<YOUR_BUCKET_PATH>/spark-sql-kafka-0-10_2.12-3.5.3.jar,oss://<YOUR_BUCKET_PATH>/spark-token-provider-kafka-0-10_2.12-3.5.3.jar
  6. 公開後、[O&M に移動] をクリックします。次に、[開始] をクリックします。

ステップ 6: 結果を表示する

Kafka からの読み取り

  1. タスクが開始したら、[タスクオーケストレーション] > [ストリーミングタスク] タブで、ターゲットタスクをクリックします。

  2. [ログ分析] タブで、ログを表示できます。

Kafka への書き込み

  1. タスクが開始したら、Kafka コンソールにログインし、インスタンスページに移動します。

  2. [メッセージクエリ] タブで、[クエリ] をクリックして、Spark が Kafka に書き込んだメッセージを表示します。

付録

以下は、pom.xml ファイルのサンプルです。

<dependencies>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.12</artifactId>
            <version>3.5.2</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.12</artifactId>
            <version>3.5.2</version>
            <scope>provided</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-hive_2.12</artifactId>
            <version>3.5.2</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>net.alchim31.maven</groupId>
                <artifactId>scala-maven-plugin</artifactId>
                <version>3.2.2</version>
                <executions>
                    <execution>
                        <goals>
                            <goal>compile</goal>
                            <goal>testCompile</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>

            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>3.3.0</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                    <archive>
                        <manifest>
                            <mainClass>com.example.Main</mainClass>
                        </manifest>
                    </archive>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>