以下示例演示如何构建连接到 R2 Data Catalog 的 Apache Spark ↗ 应用程序(使用 Scala)。该应用程序设计为在本地运行,但也可以适配为在集群上运行。
- 注册 Cloudflare 账户 ↗。
- 创建 R2 存储桶并启用数据目录。
- 创建 R2 API 令牌,并授予 R2 和数据目录权限。
- 安装 Java 17、Spark 3.5.3 和 SBT 1.10.11
- 注意:本示例中工具的特定版本对于正常运行至关重要。
- 提示:"SDKMAN" ↗ 是安装 SDK 的便捷包管理器。
首先,在计算机上创建一个新的空项目目录。
在该目录中,在 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")要启用 sbt-assembly 插件 ↗(用于构建 fat JAR),在 project/assembly.sbt 新文件中添加以下内容:
addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "1.2.0")确保 Java、Spark 和 sbt 已安装并在 shell 中可用。如果使用 SDKMAN,可以按如下方式安装:
sdk install java 17.0.14-amzn
sdk install spark 3.5.3
sdk install sbt 1.10.11安装完成后,可以使用 sbt 构建项目。这将生成一个打包的 JAR 文件。
sbt clean assembly构建完成后,输出 JAR 应位于 target/scala-2.12/R2DataCatalogDemo-assembly-1.0.jar。
要运行应用程序,需要使用 spark-submit。以下是一个示例 shell 脚本(submit.sh),包含 Spark 在 Java 17 上运行所需的 Java 兼容性标志:
# 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
- example
- com
- scala
- main
提交作业前,请确保已设置 catalog URI、warehouse 和 Cloudflare API 令牌 所需的环境变量。
export CATALOG_URI=
export WAREHOUSE=
export TOKEN=现在可以运行作业:
./submit.sh