Spark
Apache Spark 是一个被广泛使用的统一分析引擎,用于大规模数据处理。Ozone 可以作为 Spark 应用的可扩展存储层,让你使用熟悉的 Spark API 直接从 Ozone 集群读写数据。
note
本指南介绍 Apache Spark 3.x。示例已在 Spark 3.5.x 和 Apache Ozone 2.1.0 环境下测试通过。
概述
Spark 主要通过 OzoneFileSystem 连接器与 Ozone 交互,该连接器支持使用 ofs:// URI 方案进行访问。Spark 也可以通过 S3 网关使用 s3a:// 协议访问 Ozone,这有助于在不修改应用代码的情况下,将现有的云原生 Spark 应用迁移到 Ozone。
旧的 o3fs:// 方案仍为兼容性而保留,但不建议在新部署中使用。
主要优势包括:
- 将 Spark 作业生成或消费的大型数据集直接存储在 Ozone 中。
- 利用 Ozone 的可扩展性和对象存储特性来支撑 Spark 工作负载。
- 使用标准的 Spark DataFrame 和
RDDAPI 与 Ozone 数据交互。
前置条件
- Ozone 集群: 一个正在运行的 Ozone 集群。
- Ozone 客户端 JAR:
ozone-filesystem-hadoop3-client-*.jar必须位于 Spark driver 和 executor 的类路径上。 - Hadoop 3.4.x 运行时(Ozone 2.1.0+): Ozone 2.1.0 移除了若干 Hadoop 类的内置副本(
LeaseRecoverable、SafeMode、SafeModeAction),现在需要从运行时类路径中获取它们(HDDS-13574)。由于 Spark 3.5.x 自带的是 Hadoop 3.3.4,你必须将hadoop-common-3.4.x.jar与现有的 Hadoop JAR 一起添加到 Spark 类路径中。 - 配置: Spark 需要访问 Ozone 配置(
core-site.xml,以及可能的ozone-site.xml)才能连接到 Ozone 集群。
配置
1. Core Site(core-site.xml)
关于 core-site.xml 的配置,请参阅 Ozone 文件系统 (ofs) 配置章节。
2. Spark 配置(spark-defaults.conf 或 --conf)
虽然 Spark 通常会从类路径上的 core-site.xml 中读取设置,但有时仍有必要显式指定实现类:
spark.hadoop.fs.ofs.impl=org.apache.hadoop.fs.ozone.RootedOzoneFileSystem3. 安全性(Kerberos)
如果你的 Ozone 和 Spark 集群启用了 Kerberos,Spark 需要获得获取 Ozone 委托令牌的权限。
在 spark-defaults.conf 中或通过 --conf 配置以下属性,并指定你的 Ozone 文件系统 URI:
# For YARN deployments in spark3+
spark.kerberos.access.hadoopFileSystems=ofs://ozone1/将 ozone1 替换为你的 OM 服务 ID。确保运行 Spark 作业的用户拥有有效的 Kerberos 票据(kinit)。
使用示例
你可以像使用其他任何 Hadoop 兼容文件系统一样,通过 ofs:// URI 读写数据。
URI 格式: ofs://<om-service-id>/<volume>/<bucket>/path/to/key
读取数据(Scala)
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder.appName("Ozone Spark Read Example").getOrCreate()
// Read a CSV file from Ozone
val df = spark.read.format("csv")
.option("header", "true")
.option("inferSchema", "true")
.load("ofs://ozone1/volume1/bucket1/input/data.csv")
df.show()写入数据(Scala)
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder.appName("Ozone Spark Write Example").getOrCreate()
// Assume 'df' is a DataFrame you want to write
val data = Seq(("Alice", 1), ("Bob", 2), ("Charlie", 3))
val df = spark.createDataFrame(data).toDF("name", "id")
// Write DataFrame to Ozone as Parquet files
df.write.mode("overwrite")
.parquet("ofs://ozone1/volume1/bucket1/output/users.parquet")读取数据(Python)
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("Ozone Spark Read Example").getOrCreate()
# Read a CSV file from Ozone
df = spark.read.format("csv") \
.option("header", "true") \
.option("inferSchema", "true") \
.load("ofs://ozone1/volume1/bucket1/input/data.csv")
df.show()写入数据(Python)
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("Ozone Spark Write Example").getOrCreate()
# Assume 'df' is a DataFrame you want to write
data = [("Alice", 1), ("Bob", 2), ("Charlie", 3)]
columns = ["name", "id"]
df = spark.createDataFrame(data, columns)
# Write DataFrame to Ozone as Parquet files
df.write.mode("overwrite") \
.parquet("ofs://ozone1/volume1/bucket1/output/users.parquet")Spark on Kubernetes
推荐的在 Kubernetes 上使用 Ozone 运行 Spark 的方式是,将 ozone-filesystem-hadoop3-client-*.jar JAR、hadoop-common-3.4.x.jar JAR(如果使用的是 Ozone 2.1.0 及以上版本)以及 core-site.xml 直接打包进自定义 Spark 镜像中。
构建自定义 Spark 镜像
将 Ozone 客户端 JAR 和 Hadoop 兼容性 JAR 放入 /opt/spark/jars/ 目录(该目录位于 Spark 默认的 classpath 中),并将 core-site.xml 放入 /opt/spark/conf/ 目录:
FROM apache/spark:3.5.8-scala2.12-java11-python3-ubuntu
USER root
ADD https://repo1.maven.org/maven2/org/apache/ozone/ozone-filesystem-hadoop3-client/2.1.0/ozone-filesystem-hadoop3-client-2.1.0.jar \
/opt/spark/jars/
# Ozone 2.1.0+ requires Hadoop 3.4.x classes (HDDS-13574).
# Add alongside (not replacing) Spark's bundled hadoop-common-3.3.4.jar.
ADD https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-common/3.4.2/hadoop-common-3.4.2.jar \
/opt/spark/jars/
COPY core-site.xml /opt/spark/conf/core-site.xml
COPY ozone_write.py /opt/spark/work-dir/ozone_write.py
USER spark其中 core-site.xml 至少包含:
<?xml version="1.0"?>
<?xml-stylesheet type="text/xsl" href="configuration.xsl"?>
<configuration>
<property>
<name>fs.ofs.impl</name>
<value>org.apache.hadoop.fs.ozone.RootedOzoneFileSystem</value>
</property>
<property>
<name>ozone.om.address</name>
<value>om-host.example.com:9862</value>
</property>
</configuration>提交 Spark 作业
./bin/spark-submit \
--master k8s://https://YOUR_KUBERNETES_API_SERVER:6443 \
--deploy-mode cluster \
--name spark-ozone-example \
--conf spark.executor.instances=2 \
--conf spark.kubernetes.container.image=YOUR_REPO/spark-ozone:latest \
--conf spark.kubernetes.authenticate.driver.serviceAccountName=spark \
--conf spark.kubernetes.namespace=YOUR_NAMESPACE \
local:///opt/spark/work-dir/ozone_write.py将 YOUR_KUBERNETES_API_SERVER、YOUR_REPO 和 YOUR_NAMESPACE 替换为你环境中的实际值。
使用 S3A 协议
Spark 也可以通过 S3 Gateway 使用 s3a:// 协议访问 Ozone。这在将现有的云原生 Spark 应用迁移到 Ozone 而无需修改应用代码时非常有用。
有关配置细节,请参阅 S3A 文档。
评论
登录后参与评论
KnowForge