Spark访问外部存储有哪些常见方式?
Spark的核心能力之一就是灵活对接各种外部存储系统,无论你面对的是传统HDFS、云上S3,还是结构化数据库,都能通过统一的DataSource API实现读写。Spark访问外部存储的本质是借助Connector和配置驱动,让数据源与计算引擎解耦,开发者只需关注逻辑而非底层传输细节。
基于DataSource API的通用接口
Spark SQL和DataFrame从设计之初就支持通过format和option指定外部存储类型与连接参数,大多数情况下,你只需在SparkSession中调用read.format("...").load(),就能完成数据加载,这种接口的通用性体现在:
- 支持多种格式:Parquet、ORC、JSON、CSV、Avro以及JDBC。
- 通过
option传递连接信息,如路径、密钥、并发数。 - 对应的Connector包(如
spark-sql-kafka、spark-hbase)需提前引入依赖。
操作路径示例:假设你需要从远程PostgreSQL读取订单表,只需配置url、dbtable、user和password,Spark便会自动将谓词下推至数据库,减少数据传输量,对于流式场景,Structured Streaming同样支持从Kafka、文件流等源持续消费。
不同存储系统的连接配置
不同外部存储的配置差异主要体现在认证方式和协议前缀上,以下列出常见场景的关键配置项:
- HDFS:原生集成,指定
hdfs://路径即可,需确保Hadoop配置文件(core-site.xml、hdfs-site.xml)在Classpath中,或通过spark.hadoop.前缀设置。 - S3(Amazon S3及兼容对象存储):使用
s3a://协议,需添加hadoop-aws依赖,并配置spark.hadoop.fs.s3a.access.key和spark.hadoop.fs.s3a.secret.key,针对地域性部署(如华北节点),可设置spark.hadoop.fs.s3a.endpoint提升访问速度。 - HBase:通过
Spark-HBase Connector或HBaseContext操作,需指定ZooKeeper quorum和znode parent,将HBase表映射为DataFrame。 - Kafka:直接使用
spark-sql-kafka,通过subscribe指定主题,并设置kafka.bootstrap.servers,在0.10版本以上,支持Exactly-Once语义。
配置时最易踩坑的是依赖版本冲突,建议使用与Spark发行版捆绑的Connector版本,或通过--packages参数自动解析。
读写操作中的权限与性能调优
外部存储通常带有权限管控,Kerberos认证是Hadoop生态的常见门槛,你需要在spark-submit命令中指定--principal和--keytab,并在JVM参数中开启java.security.krb5.conf,对于S3,IAM角色或临时凭证(STS)是更安全的选择。
性能方面,核心约束来自网络IO和存储吞吐,常用优化手段包括:
- 使用列式格式(Parquet/ORC)并开启谓词下推,只读取需要的列和分区。
- 调整
spark.sql.files.maxPartitionBytes控制分区大小,避免小文件过多。 - 对于S3等对象存储,启用
fs.s3a.fast.upload和提升fs.s3a.threads.max来加速写入。 - 通过
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version=2提升提交效率。
Spark外部集群组件如何选择:YARN vs Kubernetes
当Spark需要借助外部集群资源管理器运行时,YARN和Kubernetes是最主流的两种选择。Spark外部集群组件的选取本质是在现有基础设施与运维成本之间做权衡,YARN适合传统大数据平台,Kubernetes则更适合云原生和容器化环境。
Spark on YARN的部署与配置
YARN作为Hadoop生态的默认资源管理器,与Spark的集成已非常成熟,部署时只需将Spark包放在所有节点,通过spark-submit --master yarn提交任务,关键参数包括:
--deploy-mode:client(提交机作为Driver)或cluster(Driver在YARN节点上运行)。spark.yarn.archive:将Spark jars打包上传至HDFS,避免每次分发。spark.yarn.maxAppAttempts:最多重试次数。
实际操作中,你需要在yarn-site.xml中配置资源分配限制,并确保HADOOP_CONF_DIR指向正确,日志查看可通过yarn logs -applicationId命令获取,对于生产环境,通常使用cluster模式保证Driver高可用。
Spark on Kubernetes的容器化实践
从Spark 2.3开始,Kubernetes成为官方支持的目标,你只需构建一个包含Spark和依赖的Docker镜像,然后通过
spark-submit --master k8s://启动,关键配置包括:
spark.kubernetes.container.image:指定镜像地址。spark.kubernetes.driver.limit.cores和spark.kubernetes.executor.limit.cores:资源限制。spark.kubernetes.authenticate.driver.serviceAccountName:服务账号。
动态资源分配在Kubernetes上需要配合spark.kubernetes.allocation.driver.node和spark.dynamicAllocation.enabled设置,相比YARN,Kubernetes的Pod启动延迟稍高,但隔离性更强,且支持自定义调度器。
两种模式的适用场景对比
| 对比维度 | YARN | Kubernetes |
|---|---|---|
| 部署复杂度 | 依赖Hadoop集群,需提前搭建HDFS和YARN | 需要容器编排环境,镜像构建有学习成本 |
| 资源隔离 | 基于CGroup,但隔离粒度较粗 | 基于CRI和Namespace,CPU和内存隔离更严格 |
| 弹性伸缩 | 依赖YARN资源申请机制,扩缩较慢 | 利用Horizontal Pod Autoscaler,响应更快 |
| 成本 | 通常使用裸金属或虚拟机,资源利用率中等 | 可结合竞价实例,按需付费,适合云上场景 |
行业共识认为,如果已有Hadoop大数据平台,YARN的兼容性和稳定性更优;若团队正在转向微服务或云原生架构,Kubernetes能统一资源层,减少重复维护。
Spark读写外部数据源性能优化技巧
无论外部存储是HDFS还是S3,IO性能瓶颈往往决定作业效率。Spark读写外部数据源时,优化方向集中在减少数据扫描量、提升并行度以及利用缓存策略。
减少数据读取的IO开销
- 优先使用列式存储格式(Parquet/ORC),并开启
spark.sql.parquet.pushDownPredicate和spark.sql.parquet.filterPushdown,让Spark只读取符合条件的行与列。 - 对分区表按分区键过滤,避免全表扫描,在
where子句限定日期范围,Spark会直接跳过无关分区目录。 - 利用数据本地性(Data Locality),通过
参数协调等待时间,尽量让计算任务在数据所在节点运行。spark.locality.wait
合理设置并行度与连接参数
- 调整
spark.sql.files.maxPartitionBytes(默认128MB)控制每个分区读取的数据量,避免小文件导致过多任务。 - 对于对象存储(如S3),增大
fs.s3a.connection.maximum和fs.s3a.threads.max,提升并发连接数。 - 为外部数据库设置
fetchSize和batchSize,例如JDBC连接时指定batchsize=1000,减少网络往返。
缓存与广播变量的使用
- 如果同一外部数据被多次使用,可以调用
cache()或persist()存入内存(或磁盘),后续操作直接读取缓存。 - 对于维度表等小数据集,使用
broadcast广播变量,避免Shuffle,同时加速Join操作。 - 在流式任务中,合理设置
spark.sql.streaming.fileSink.log.deletion,清理过期日志,避免元数据膨胀。
Spark访问外部存储与集群组件常见问题
问题1:Spark访问HDFS时需要额外配置吗?
如果Spark与HDFS部署在同一集群,通常只需将core-site.xml和hdfs-site.xml复制到Spark的conf目录下,或通过spark.hadoop.前缀设置,若跨集群访问,则需确保两个集群的Kerberos互通,并在spark-submit中指定Principal和Keytab。
问题2:Spark on Kubernetes如何配置动态资源分配?
在Kubernetes模式下,动态资源分配需要手动开启spark.dynamicAllocation.enabled=true,并设置spark.dynamicAllocation.maxExecutors限制上限,需要配置spark.kubernetes.allocation.driver.node以匹配Pod调度器,注意,Kubernetes的Executor创建有秒级延迟,短时间任务可能不适合开启此功能。
问题3:从S3读取数据比从HDFS慢,如何优化?
S3属于对象存储,其延迟通常高于HDFS,优化方向包括:使用S3A协议并开启fast.upload和direct模式;增大fs.s3a.readahead.range预读范围;将数据转为Parquet格式并压缩;若读取频繁,可考虑使用S3的缓存层(如S3A与本地磁盘的混合模式)。
首发原创文章,作者:王坚,如若转载,请注明出处:https://idctop.com/article/542215.html



