集成

Spark

师成师成· 更新于 2026-09-28· 阅读 12 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

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 和 RDD API 与 Ozone 数据交互。

前置条件

  1. Ozone 集群: 一个正在运行的 Ozone 集群。
  2. Ozone 客户端 JAR: ozone-filesystem-hadoop3-client-*.jar 必须位于 Spark driver 和 executor 的类路径上。
  3. 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 类路径中。
  4. 配置: 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.RootedOzoneFileSystem

3. 安全性(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 文档。

评论

登录后参与评论

正在加载评论…