このトピックでは、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 パッケージが含まれています。このステップはスキップできます。 別のエンジンバージョンを使用する場合は、次のステップを実行する必要があります。
Spark のバージョンに基づいて、次の依存関係を手動で追加できます。
Spark 3.5: kafka-spark35-jars.zip。
Spark 3.4: kafka-spark34-jars.zip。
Spark 3.3: kafka-spark33-jars.zip。
ファイルを解凍した後、すべての 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 パッケージをアップロードする
ファイルアップロードページに移動します。
E-MapReduce コンソールにログインします。
左側のナビゲーションウィンドウで、 を選択します。
[Spark] ページで、ターゲットワークスペースをクリックします。
EMR Serverless Spark ページで、左側のナビゲーションウィンドウの [ファイル管理] をクリックします。
[ファイル管理] ページで、[ファイルのアップロード] をクリックします。
[ファイルのアップロード] ダイアログボックスで、アップロードエリアをクリックするか、パッケージをエリアにドラッグして JAR パッケージを選択できます。
ステップ 5: ストリーミングタスクを作成して実行する
EMR Serverless Spark ページで、左側のナビゲーションウィンドウの [データ開発] をクリックします。
[開発] タブで、新規アイコン
をクリックします。表示されるダイアログボックスで、[名前] を入力し、[ストリーミングタスク] の JAR タイプを選択して、[OK] をクリックします。
右上隅で、ターゲットキューを選択します。
キューの追加方法の詳細については、「リソースキューの管理」をご参照ください。
新しい開発タブで、次のパラメーターを構成し、他のパラメーターはデフォルト設定のままにします。次に、[公開] をクリックします。
パラメーター
説明
メイン 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公開後、[O&M に移動] をクリックします。次に、[開始] をクリックします。
ステップ 6: 結果を表示する
Kafka からの読み取り
タスクが開始したら、 タブで、ターゲットタスクをクリックします。
[ログ分析] タブで、ログを表示できます。
Kafka への書き込み
タスクが開始したら、Kafka コンソールにログインし、インスタンスページに移動します。
[メッセージクエリ] タブで、[クエリ] をクリックして、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>