Spark 3.5.1无法读取基于Ceph的本地S3存储对象内容求助
问题:Spark 3.5.1读取Ceph S3存储文件能获取元数据但无法读取内容
使用Spark 3.5.1连接基于Ceph的本地S3存储,已成功访问存储桶、列出文件,且能获取到文件的大小(DEBUG日志显示s3a://input/testfile.csv: testfile.csv size=18),但读取文件内容时返回空DataFrame。
环境与代码
PySpark代码如下:
from pyspark import SparkConf from pyspark.sql.types import * from pyspark.sql import SparkSession conf = SparkConf() conf.setAll([ ("spark.jars.packages", "org.apache.hadoop:hadoop-aws:3.3.6"), ("spark.hadoop.fs.s3a.access.key", "R*************6"), ("spark.hadoop.fs.s3a.secret.key", "1***************e"), ("spark.hadoop.fs.s3a.endpoint", "192.168.52.63:8000") ]) spark = SparkSession.builder \ .master("local[*]") \ .config(conf=conf) \ .appName("read-s3-with-spark") \ .getOrCreate() schema = StructType( [StructField('name', StringType(), True), StructField('int1', IntegerType(), True), StructField('int2', IntegerType(), True) ] ) spark.sparkContext.setLogLevel("DEBUG") df = spark.read \ .option("header", "true") \ .schema(schema) \ .csv("s3a://input/testfile.csv", sep=' ') df.show(n=1)
日志信息
DEBUG日志片段:
DEBUG S3AFileSystem: s3a://input/testfile.csv: testfile.csv size=18
INFO日志片段:
24/05/20 02:35:00 INFO MetricsSystemImpl: s3a-file-system metrics system started 24/05/20 02:35:01 INFO MetadataLogFileIndex: Reading streaming file log from s3a://input/testfile.csv/_spark_metadata 24/05/20 02:35:01 INFO FileStreamSinkLog: BatchIds found from listing: 24/05/20 02:35:03 INFO FileSourceStrategy: Pushed Filters: 24/05/20 02:35:03 INFO FileSourceStrategy: Post-Scan Filters: 24/05/20 02:35:03 INFO CodeGenerator: Code generated in 176.139675 ms 24/05/20 02:35:03 INFO MemoryStore: Block broadcast_0 stored as values in memory (estimated size 496.6 KiB, free 4.1 GiB) 24/05/20 02:35:03 INFO MemoryStore: Block broadcast_0_piece0 stored as bytes in memory (estimated size 54.4 KiB, free 4.1 GiB) 24/05/20 02:35:03 INFO BlockManagerInfo: Added broadcast_0_piece0 in memory on master:38197 (size: 54.4 KiB, free: 4.1 GiB) 24/05/20 02:35:03 INFO SparkContext: Created broadcast 0 from showString at NativeMethodAccessorImpl.java:0 24/05/20 02:35:03 INFO FileSourceScanExec: Planning scan with bin packing, max size: 4194304 bytes, open cost is considered as scanning 4194304 bytes. +----+----+----+ |name|int1|int2| +----+----+----+ +----+----+----+ 24/05/20 02:35:04 INFO SparkContext: Invoking stop() from shutdown hook 24/05/20 02:35:04 INFO SparkContext: SparkContext is stopping with exitCode 0
目标文件内容
testfile.csv的实际内容:
name int1 int2 first 1 2 second 3 4
问题排查与解决建议
补充S3A核心配置参数
本地Ceph S3集群需要启用路径样式访问,并禁用SSL(无证书场景),在SparkConf中添加以下配置:("spark.hadoop.fs.s3a.path.style.access", "true"), ("spark.hadoop.fs.s3a.connection.ssl.enabled", "false"),验证文件完整性
DEBUG日志显示文件大小仅18字节,但你提供的文件实际内容远大于该尺寸,说明Spark从Ceph获取的文件内容不完整。请:- 检查Ceph存储中
testfile.csv的实际大小,确认文件是否完整上传 - 重新上传文件至Ceph,排除上传过程中的截断问题
- 检查Ceph存储中
调整分隔符配置
若文件中字段为多个空格分隔,使用sep=' '只会匹配单个空格,导致字段解析失败。将分隔符改为匹配任意空白字符:.option("sep", "\\s+")对齐Hadoop-AWS版本
Spark 3.5.1默认依赖Hadoop 3.3.4,当前使用的hadoop-aws:3.3.6存在小版本差异,可尝试替换为匹配版本:("spark.jars.packages", "org.apache.hadoop:hadoop-aws:3.3.4"),
内容的提问来源于stack exchange,提问作者Amin mosayyebzadeh
相关产品推荐
相关产品推荐

