Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

Spark (Scala)

最終更新 Markdown で表示Agent セットアップ

Apache Spark アプリケーション(Scala)を R2 Data Catalog に接続する例です。このアプリケーションはローカル実行向けですが、クラスターでも動くように応用できます。

前提条件

使用例

まず、マシン上の適当な場所に空のプロジェクトディレクトリを作成します。

そのディレクトリ内に、src/main/scala/com/example/R2DataCatalogDemo.scala を作成します。これが Spark アプリケーションのメインエントリポイントになります。

package com.example

import org.apache.spark.sql.SparkSession

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

        val uri = sys.env("CATALOG_URI")
        val warehouse = sys.env("WAREHOUSE")
        val token = sys.env("TOKEN")

        val spark = SparkSession.builder()
            .appName("My R2 Data Catalog Demo")
            .master("local[*]")
            .config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
            .config("spark.sql.catalog.mydemo", "org.apache.iceberg.spark.SparkCatalog")
            .config("spark.sql.catalog.mydemo.type", "rest")
            .config("spark.sql.catalog.mydemo.uri", uri)
            .config("spark.sql.catalog.mydemo.warehouse", warehouse)
            .config("spark.sql.catalog.mydemo.token", token)
            .getOrCreate()

        import spark.implicits._

        val data = Seq(
            (1, "Alice", 25),
            (2, "Bob", 30),
            (3, "Charlie", 35),
            (4, "Diana", 40)
        ).toDF("id", "name", "age")

        spark.sql("USE mydemo")

        spark.sql("CREATE NAMESPACE IF NOT EXISTS demoNamespace")

        data.writeTo("demoNamespace.demotable").createOrReplace()

        val readResult = spark.sql("SELECT * FROM demoNamespace.demotable WHERE age > 30")
        println("Records with age > 30:")
        readResult.show()
    }
}

このアプリケーションのビルドと依存関係の管理には sbt(simple build tool) を使います。次は、プロジェクトのルートに置く build.sbt の例です。必要な依存関係をまとめた fat JAR を作る設定です。

name := "R2DataCatalogDemo"

version := "1.0"

val sparkVersion = "3.5.3"
val icebergVersion = "1.8.1"

// You need to use binaries of Spark compiled with either 2.12 or 2.13; and 2.12 is more common.
// If you download Spark 3.5.3 with sdkman, then it comes with 2.12.18
scalaVersion := "2.12.18"

libraryDependencies ++= Seq(
    "org.apache.spark" %% "spark-core" % sparkVersion,
    "org.apache.spark" %% "spark-sql" % sparkVersion,
    "org.apache.iceberg" % "iceberg-core" % icebergVersion,
    "org.apache.iceberg" % "iceberg-spark-runtime-3.5_2.12" % icebergVersion,
    "org.apache.iceberg" % "iceberg-aws-bundle" % icebergVersion,
)

// build a fat JAR with all dependencies
assembly / assemblyMergeStrategy := {
    case PathList("META-INF", "services", xs @ _*) => MergeStrategy.concat
    case PathList("META-INF", xs @ _*) => MergeStrategy.discard
    case "reference.conf" => MergeStrategy.concat
    case "application.conf" => MergeStrategy.concat
    case x if x.endsWith(".properties") => MergeStrategy.first
    case x => MergeStrategy.first
}

// For Java  17 Compatibility
Compile / javacOptions ++= Seq("--release", "17")

fat JAR のビルドに使う sbt-assembly プラグイン を有効にするには、project/assembly.sbt に次を追加します。

addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "1.2.0")

Java、Spark、sbt がインストールされ、シェルから使えることを確認します。SDKMAN を使う場合は、次のようにインストールできます。

sdk install java 17.0.14-amzn
sdk install spark 3.5.3
sdk install sbt 1.10.11

すべてインストールしたら、sbt でプロジェクトをビルドします。依存関係をまとめた 1 つの JAR が生成されます。

sbt clean assembly

ビルド後の JAR は target/scala-2.12/R2DataCatalogDemo-assembly-1.0.jar にあります。

アプリケーションの実行には spark-submit を使います。次は、Java 17 上の Spark に必要な互換性フラグを含むシェルスクリプト(submit.sh)の例です。

# We need to set these "--add-opens" so that Spark can run on Java 17 (it needs access to
# parts of the JVM which have been modularized and made internal).
JAVA_17_COMPATIBILITY="--add-opens=java.base/sun.nio.ch=ALL-UNNAMED --add-opens=java.base/java.nio=ALL-UNNAMED --add-opens=java.base/java.lang=ALL-UNNAMED --add-opens=java.base/java.util=ALL-UNNAMED --add-opens=java.base/java.util.concurrent=ALL-UNNAMED"

spark-submit \
--conf "spark.driver.extraJavaOptions=$JAVA_17_COMPATIBILITY" \
--conf "spark.executor.extraJavaOptions=$JAVA_17_COMPATIBILITY" \
--class com.example.R2DataCatalogDemo target/scala-2.12/R2DataCatalogDemo-assembly-1.0.jar

実行する前に、スクリプトに実行権限を付けます。

chmod +x submit.sh

この時点で、プロジェクトディレクトリは次のような構成になります。

  • Makefile
  • README.md
  • build.sbt
  • project
    • assembly.sbt
    • build.properties
    • project
  • spark-submit.sh
  • src
    • main
      • scala
        • com
          • example
            • R2DataCatalogDemo.scala

ジョブを投入する前に、カタログ URI、warehouse、Cloudflare API トークン に必要な環境変数を設定します。

export CATALOG_URI=
export WAREHOUSE=
export TOKEN=

これでジョブを実行できます。

./submit.sh

役に立ちましたか?