首页
/ MinIO 大数据实战:分离式 HDP Spark 与 Hive 架构下的 S3A 调优全指南

MinIO 大数据实战:分离式 HDP Spark 与 Hive 架构下的 S3A 调优全指南

2026-09-05 15:52:42作者:薛曦旖Francesca

本文围绕 MinIO 官方文档《Disaggregated HDP Spark and Hive with MinIO》,讲解如何在 Hortonworks 发行版(HDP)环境中,将 Hadoop、Spark、Hive 的存储层从 HDFS 迁移到 MinIO 对象存储:包括 S3A 连接器(s3a scheme)的核心参数、S3A 目录暂存提交器(directory staging committer)原理与写放大规避,以及用 Spark Pi、WordCount 验证存储链路的完整实操。读完后你能够独立完成一套“计算无状态、存储对象化”的云原生大数据集群配置,并理解 MinIO 侧支撑这些负载的 S3 / S3 SELECT API 实现。

MinIO 云原生大数据架构

1. 云原生分离式架构:计算与存储解耦

文档给出的参考架构(如上架构图)遵循云原生大数据的经典分层:

  • 计算层(无状态):Kubernetes 在计算节点上弹性地管理 Spark 与 Hive 的无状态容器。Spark 与 Kubernetes 有原生的调度器集成;Hive 因历史原因在 Kubernetes 之上继续使用 YARN 调度。
  • 存储层(有状态):MinIO 容器同样由 Kubernetes 管理,但作为有状态容器运行,本地存储(JBOD/JBOF)以持久化本地卷(persistent local volumes)的方式挂载,构成纠删码分布式对象存储。这种部署方式天然支持多租户:不同客户的数据在桶/密钥层面相互隔离。
  • 访问路径:所有对 MinIO 的访问都通过 S3 API 与 S3 SELECT(SQL SELECT)API,即 Hadoop 生态标准连接器 s3a 可直接对接,无需自研文件系统。
  • 多集群/多站点联邦:MinIO 支持类似 AWS region 与存储分层(tiers)的多集群、多站点联邦。利用 MinIO 的 ILM(Information Lifecycle Management,信息生命周期管理),可将数据在 NVMe 热存储与 HDD 温存储之间自动分层。
  • 安全:全部数据使用每对象独立密钥加密;跨租户的访问控制与身份管理由 MinIO 统一承担,支持 OpenID Connect 或 Kerberos/LDAP/AD。

这套架构的核心收益是:计算节点可完全无状态化并弹性伸缩,存储容量独立扩展,且通过 S3 语义获得与公有云一致的移植性。

2. 前置条件(Prerequisites)

在开始配置前,需要准备好两套组件:

  1. Hortonworks 发行版(HDP)
    • 安装 Ambari,它会自动完成 YARN 的安装与配置;
    • 通过 HDP 安装包安装 Spark(文档基于 HDP 3.0.1 时代的安装指引)。
  2. MinIO 分布式服务器,二选一部署:
    • 基于 Kubernetes 的部署方式;
    • 基于 MinIO Helm Chart 的部署方式。本仓库根目录的 index.yamlhelm-releases 目录提供了历次发布的 MinIO Helm Chart 存档(minio-x.y.z.tgz),可从中选取与集群兼容的 chart 版本离线安装。

完成安装后,进入 Ambari UI http://<ambari-server>:8080/,使用默认凭据(用户名:admin,密码:admin)登录,后续所有组件配置都通过该界面完成。

3. 配置 Hadoop:S3A 连接器与提交器调优

在 Ambari 中按 Services -> HDFS -> CONFIGS -> ADVANCED 路径进入 HDFS 的高级配置区,然后打开 Custom core-site,这里是配置 s3a 连接 MinIO 的核心位置。

3.1 配置修改工作流:kv-pairify 技巧

core-site.xml 是 XML 格式,直接编辑易错。文档给出了一套基于 yq + jq 的转换工具链,把配置转成 name=value 的键值对形式,方便查看与检索:

sudo pip install yq
alias kv-pairify='yq ".configuration[]" | jq ".[]" | jq -r ".name + \"=\" + .value"'

使用它可以从 core-site.xml 中快速过滤出某类配置,例如查看 MapReduce 相关项:

cat ${HADOOP_CONF_DIR}/core-site.xml | kv-pairify | grep "mapred"

3.2 MapReduce / 任务级调优参数

文档以一个 12 台计算节点、内存总量约 1.2 TiB 的集群为基准,给出下列针对大内存节点的 MapReduce 优化项,写入 core-site.xml

mapred.maxthreads.generate.mapoutput=2 # Num threads to write map outputs
mapred.maxthreads.partition.closer=0 # Asynchronous map flushers
mapreduce.fileoutputcommitter.algorithm.version=2 # Use the latest committer version
mapreduce.job.reduce.slowstart.completedmaps=0.99 # 99% map, then reduce
mapreduce.reduce.shuffle.input.buffer.percent=0.9 # Min % buffer in RAM
mapreduce.reduce.shuffle.merge.percent=0.9 # Minimum % merges in RAM
mapreduce.reduce.speculative=false # Disable speculation for reducing
mapreduce.task.io.sort.factor=999 # Threshold before writing to disk
mapreduce.task.sort.spill.percent=0.9 # Minimum % before spilling to disk

这些参数的共同意图是:把 shuffle 尽量留在内存中完成(99% map 完成才启动 reduce、0.9 的 buffer/merge 比例、999 的排序因子相当于几乎不提前落盘),并关闭 reduce 端推测执行以减少对象存储上的重复写入。

3.3 S3A 提交器原理:为什么要用 directory staging committer

这是本文档最有原理价值的部分,决定了对象存储上 MapReduce 写入性能的上限:

  • 背景:S3A 是 Hadoop 访问 S3 及 MinIO 等 S3 兼容存储的连接器。MapReduce 作业与对象存储的交互方式与 HDFS 相同,而它们依赖 HDFS 的原子 rename 来完成“写数据—提交”语义。对象存储的操作本身是原子的,但不提供 rename API
  • 默认 committer 的代价:默认 S3A committer 通过 copy + delete 两个操作来模拟 rename。这种交互模式带来严重的写放大(write amplification);
  • Netflix 的解法:Netflix 开发了两个新的暂存提交器(staging committer)——Directory staging committerPartitioned staging committer,两者都不需要 rename 操作,直接利用对象存储的原生语义;
  • 结论:在三类提交器(含 Magic committer)的对比评估中,directory staging committer 性能最佳,因此下文配置将其指定为 fs.s3a.committer.name=directory

3.4 S3A 连接 MinIO 的完整参数

将以下条目加入 core-site.xml(即 Custom core-site),其中最重要的选项是 endpoint、凭证、path-style 访问与提交器选择:

fs.s3a.access.key=minio
fs.s3a.secret.key=minio123
fs.s3a.path.style.access=true
fs.s3a.block.size=512M
fs.s3a.buffer.dir=${hadoop.tmp.dir}/s3a
fs.s3a.committer.magic.enabled=false
fs.s3a.committer.name=directory
fs.s3a.committer.staging.abort.pending.uploads=true
fs.s3a.committer.staging.conflict-mode=append
fs.s3a.committer.staging.tmp.path=/tmp/staging
fs.s3a.committer.staging.unique-filenames=true
fs.s3a.connection.establish.timeout=5000
fs.s3a.connection.ssl.enabled=false
fs.s3a.connection.timeout=200000
fs.s3a.endpoint=http://minio:9000
fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem

在基础连接参数之外,文档针对大规模集群追加了一组并发与分片上传参数(这是吞吐量的关键):

fs.s3a.committer.threads=2048 # Number of threads writing to MinIO
fs.s3a.connection.maximum=8192 # Maximum number of concurrent conns
fs.s3a.fast.upload.active.blocks=2048 # Number of parallel uploads
fs.s3a.fast.upload.buffer=disk # Use disk as the buffer for uploads
fs.s3a.fast.upload=true # Turn on fast upload mode
fs.s3a.max.total.tasks=2048 # Maximum number of parallel tasks
fs.s3a.multipart.size=512M # Size of each multipart chunk
fs.s3a.multipart.threshold=512M # Size before using multipart uploads
fs.s3a.socket.recv.buffer=65536 # Read socket buffer hint
fs.s3a.socket.send.buffer=65536 # Write socket buffer hint
fs.s3a.threads.max=2048 # Maximum number of threads for S3A

参数解读(结合对象存储上传模型):

  • fs.s3a.endpoint:指向 MinIO 服务地址,示例中为 Kubernetes 内网域名 minio:9000fs.s3a.path.style.access=true 使用路径风格寻址(http://host:9000/bucket/key),这是 MinIO 部署下的标准做法;
  • fs.s3a.block.size=512Mfs.s3a.multipart.size=512M 保持一致,块大小即分片大小,减少 multipart 分片数量;
  • fs.s3a.fast.upload=true + fs.s3a.fast.upload.buffer=disk:启用快速上传模式,上传缓冲落在本地磁盘,适合大吞吐场景;
  • 2048 级别的线程/任务并发与 8192 最大连接数,是为 1.2 TiB 内存的 12 节点集群标定的值,较小集群应按比例调低,避免压垮对象存储网关或产生大量小连接;
  • 其余 S3A 优化项可参考 Hadoop 官方 hadoop-aws 文档与 committers 文档(Hadoop 官方文档站对应条目)。

应用配置后,回到 Ambari 重启 Hadoop 相关服务使配置生效。

4. 配置 Spark2:把 S3A 参数注入 spark-defaults.conf

进入 Services -> Spark2 -> CONFIGS,打开 Custom spark-defaults,把同一套 S3A 参数以 spark.hadoop.* 前缀写入 spark-defaults.conf(Spark 会将其透传给 Hadoop FileSystem 层):

spark.hadoop.fs.s3a.access.key minio
spark.hadoop.fs.s3a.secret.key minio123
spark.hadoop.fs.s3a.path.style.access true
spark.hadoop.fs.s3a.block.size 512M
spark.hadoop.fs.s3a.buffer.dir ${hadoop.tmp.dir}/s3a
spark.hadoop.fs.s3a.committer.magic.enabled false
spark.hadoop.fs.s3a.committer.name directory
spark.hadoop.fs.s3a.committer.staging.abort.pending.uploads true
spark.hadoop.fs.s3a.committer.staging.conflict-mode append
spark.hadoop.fs.s3a.committer.staging.tmp.path /tmp/staging
spark.hadoop.fs.s3a.committer.staging.unique-filenames true
spark.hadoop.fs.s3a.committer.threads 2048 # number of threads writing to MinIO
spark.hadoop.fs.s3a.connection.establish.timeout 5000
spark.hadoop.fs.s3a.connection.maximum 8192 # maximum number of concurrent conns
spark.hadoop.fs.s3a.connection.ssl.enabled false
spark.hadoop.fs.s3a.connection.timeout 200000
spark.hadoop.fs.s3a.endpoint http://minio:9000
spark.hadoop.fs.s3a.fast.upload.active.blocks 2048 # number of parallel uploads
spark.hadoop.fs.s3a.fast.upload.buffer disk # use disk as the buffer for uploads
spark.hadoop.fs.s3a.fast.upload true # turn on fast upload mode
spark.hadoop.fs.s3a.impl org.apache.hadoop.spark.hadoop.fs.s3a.S3AFileSystem
spark.hadoop.fs.s3a.max.total.tasks 2048 # maximum number of parallel tasks
spark.hadoop.fs.s3a.multipart.size 512M # size of each multipart chunk
spark.hadoop.fs.s3a.multipart.threshold 512M # size before using multipart uploads
spark.hadoop.fs.s3a.socket.recv.buffer 65536 # read socket buffer hint
spark.hadoop.fs.s3a.socket.send.buffer 65536 # write socket buffer hint
spark.hadoop.fs.s3a.threads.max 2048 # maximum number of threads for S3A

注意:原文档中 spark.hadoop.fs.s3a.impl 一行的值写作 org.apache.hadoop.spark.hadoop.fs.s3a.S3AFileSystem,多了一段 spark.,属于原文笔误;实际 FQCN 应为 org.apache.hadoop.fs.s3a.S3AFileSystem(与 core-site.xml 中 fs.s3a.impl 的取值一致)。

应用后重启 Spark 服务。

5. 配置 Hive:hive-site.xml 的并发参数

进入 Services -> Hive -> CONFIGS -> ADVANCED,打开 Custom hive-site,向 hive-site.xml 添加下列条目。这些参数主要提升 Hive 在对象存储上的目录列举、动态分区写入与 metastore 并发能力:

hive.blobstore.use.blobstore.as.scratchdir=true
hive.exec.input.listing.max.threads=50
hive.load.dynamic.partitions.thread=25
hive.metastore.fshandler.threads=50
hive.mv.files.threads=40
mapreduce.input.fileinputformat.list-status.num-threads=50

各项含义:

  • hive.blobstore.use.blobstore.as.scratchdir=true:使用对象存储作为 scratchdir;
  • hive.exec.input.listing.max.threads:输入列举的最大线程数(对象存储上的 List 操作是常见瓶颈,多线程列举可显著加速);
  • hive.load.dynamic.partitions.thread:动态分区插入的并行写入线程;
  • hive.metastore.fshandler.threads:metastore 的 FileSystem handler 线程池;
  • hive.mv.files.threads:物化视图文件操作线程;
  • mapreduce.input.fileinputformat.list-status.num-threads:MapReduce FileInputFormat 列举文件状态的线程数。

应用后重启全部 Hive 服务。

6. 样例应用验证:Spark Pi 与 WordCount

配置完成后,用两个经典样例验证整条“Spark -> S3A -> MinIO”链路。

6.1 Spark Pi:计算密集型验证

Spark Pi 通过“投点法”估算圆周率:在单位正方形 (0,0)~(1,1) 内生成随机点,统计落入单位圆的点比例来近似 π。步骤:

  1. spark 用户登录任一装有 Spark client 的节点:
cd /usr/hdp/current/spark2-client
su spark
  1. yarn-client 模式提交 Pi 作业:
./bin/spark-submit --class org.apache.spark.examples.SparkPi \
    --master yarn-client \
    --num-executors 1 \
    --driver-memory 512m \
    --executor-memory 512m \
    --executor-cores 1 \
    examples/jars/spark-examples*.jar 10
  1. 作业输出应包含近似 π 值,例如:
17/03/22 23:21:10 INFO DAGScheduler: Job 0 finished: reduce at SparkPi.scala:38, took 1.302805 s
Pi is roughly 3.1445191445191445

作业状态也可通过 YARN ResourceManager Web UI 中的 Job History 查看。

6.2 WordCount:直接读写 MinIO 对象桶

WordCount 统计文本文件中单词出现次数,构建 (String, Int) 数据集并保存到文件。此例验证 s3a:// 为路径直接读写 MinIO 桶 的端到端能力。

步骤 1:上传输入文件到 MinIO 桶(经 S3A)

log4j.properties 为任意文本输入,hadoop fs 走的是已配置好的 S3A 文件系统:

hadoop fs -copyFromLocal /etc/hadoop/conf/log4j.properties \
          s3a://testbucket/testdata

步骤 2:启动 Spark Shell

cd /usr/hdp/current/spark2-client
su spark
./bin/spark-shell --master yarn-client --driver-memory 512m --executor-memory 512m

启动成功后会看到 Spark context / Spark session 初始化信息与 scala> 提示符(示例环境输出节选):

Spark context Web UI available at http://172.26.236.247:4041
Spark context available as 'sc' (master = yarn, app id = application_1490217230866_0002).
Spark session available as 'spark'.
...
scala>

步骤 3:在 scala> 提示符提交 WordCount 作业(请替换为你的实际桶名与路径):

scala> val file = sc.textFile("s3a://testbucket/testdata")
file: org.apache.spark.rdd.RDD[String] = s3a://testbucket/testdata MapPartitionsRDD[1] at textFile at <console>:24

scala> val counts = file.flatMap(line => line.split(" ")).map(word => (word, 1)).reduceByKey(_ + _)
counts: org.apache.spark.rdd.RDD[(String, Int)] = ShuffledRDD[4] at reduceByKey at <console>:25

scala> counts.saveAsTextFile("s3a://testbucket/wordcount")

步骤 4:查看结果

在 Scala shell 内计数:

scala> counts.count()
364

退出 shell 后用 hadoop fs 列出 MinIO 中的输出对象,结果应类似:

Found 3 items
-rw-rw-rw-   1 spark spark          0 2019-05-04 01:36 s3a://testbucket/wordcount/_SUCCESS
-rw-rw-rw-   1 spark spark       4956 2019-05-04 01:36 s3a://testbucket/wordcount/part-00000
-rw-rw-rw-   1 spark spark       5616 2019-05-04 01:36 s3a://testbucket/wordcount/part-00001

_SUCCESS 标记文件与 part-* 分片正是第 3.3 节所述 committer 提交语义的落地产物:directory staging committer 先写入暂存路径,提交后呈现为最终路径下的分片对象。

7. MinIO 侧实现视角:文档能力在源码中的对应

文档反复强调“所有对 MinIO 的访问均通过 S3 / SQL SELECT API”,这并非一句口号,可以在本仓库源码中逐一对应:

  • S3 SELECT(SQL SELECT)API:路由注册于 S3 API 路由表GET Object?select 请求经中间件后交给 SelectObjectContentHandler;该 handler 的实现入口在 S3 Select 处理器(注释即 SelectObjectContentHandler - GET Object?select),用于在对象上直接执行 SQL 谓词过滤,免去大数据作业全量拉取对象。其查询解析与执行引擎位于 internal/s3select 包(约 70 个源文件,含 CSV/JSON/Parquet 解析),支撑 Hadoop 生态直接对 MinIO 对象做 SQL 下推。
  • 多租户身份与访问控制:文档提到 MinIO 使用 OpenID Connect 或 Kerberos/LDAP/AD 管理租户间身份。仓库中 cmd/iam.gocmd/iam-object-store.gocmd/iam-etcd-store.go 构成 IAM 主体,internal/jwt/ 处理 JWT 凭据,internal/crypto/cmd/kms-handlers.go 则对应“每对象密钥加密”与可选的 KMS 后端。
  • ILM 热/温分层:文档所述“NVMe 热存储与 HDD 温存储之间分层”由 tier 子系统实现,参见 分层存储主实现分层清理器,以及各存储后端适配器 warm-backend-s3.gowarm-backend-azure.gowarm-backend-gcs.gowarm-backend-minio.go
  • multipart 上传模型:第 3.4 节的 fs.s3a.multipart.sizefs.s3a.fast.upload.active.blocks 等参数最终落在 S3 multipart 协议上,MinIO 侧的多部分上传对象管理实现于 erasure-multipart.goobject-multipart-handlers.go。这也是为什么“分片大小 = 块大小(512M)+ 高并发分片上传”能在对象存储上取得接近顺序写盘的吞吐。
  • S3 SELECT 配套示例:仓库还提供了一个 S3Zip/S3 Select 扩展示例目录 docs/extensions/s3zip,可帮助理解 S3 数据面扩展的集成方式。

8. 落地要点与参数适用边界

  1. 并发参数是基准值而非万能值:文档中的 2048 线程 / 8192 连接 / 512M 块大小是为“12 节点、1.2 TiB 内存”集群标定的。更小集群应等比例下调 fs.s3a.threads.maxfs.s3a.connection.maximumfs.s3a.max.total.tasks,否则客户端资源反而成为瓶颈;
  2. 提交器选择不可省略fs.s3a.committer.name=directory 是对象存储上 MapReduce 写入性能的关键开关;若保留默认 committer,copy+delete 模拟 rename 的写放大会抵消所有其他优化;
  3. spark.hadoop.* 前缀的一致性:Spark 侧参数与 core-site.xml 一一对应(仅加 spark.hadoop. 前缀),排查 Spark 无法访问 MinIO 的问题时,先核对两份配置是否同源(本文第 4 节还指出原文档中 fs.s3a.impl 的笔误);
  4. 安全参数:示例中 fs.s3a.connection.ssl.enabled=false 与明文 minio/minio123 凭据仅为内网演示配置,生产环境应启用 TLS 并接入 KMS/IDP;
  5. 验证顺序:建议按 Pi(计算链路)→ WordCount(对象存储读写链路)→ 真实数据集(Hive 查询)的顺序逐层验证,任何一层失败都能快速定位是 Spark/YARN 配置问题还是 S3A/凭证问题。

完成以上配置与验证后,Hadoop 生态(HDFS 语义层)即可整体运行在 MinIO 对象存储之上,计算容器保持无状态、弹性伸缩,存储通过 S3 语义独立扩展,并按 ILM 策略自动冷热分层。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.79 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
988
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384