如何使用boto3读取S3中的parquet文件并加载为Spark DataFrame
问题原因说明
报错的核心诱因有两点:
bucket.objects.filter(prefix=/my/path)返回的是S3对象集合迭代器,并非单个S3对象,直接调用get()会拿到集合类结果而非单文件字节流spark.read.parquet()仅支持传入文件路径字符串作为参数,不识别Python的BytesIO对象,因此触发类型转换异常。csv可读取是因为csv读取接口兼容文件对象传参,parquet底层实现不支持该传参逻辑。
解决代码
提前确保环境安装了pyarrow依赖,且Spark开启了Arrow优化。
1. 读取单个Parquet文件
# 开启Spark Arrow支持,提升转换效率 spark.conf.set("spark.sql.execution.arrow.enabled", "true") import io import pyarrow.parquet as pq import pyarrow as pa # 初始化S3资源 res = autorefresh_session.resource('s3') # 直接指定单个文件的完整key obj = res.Object(bucket_name='mybucket', key='my/path/myfile1.parquet') # 读取文件字节流 body = io.BytesIO(obj.get()['Body'].read()) # 用pyarrow读取parquet字节流为Arrow表 pa_table = pq.read_table(body) # 直接转换为Spark DataFrame,无需经过pandas,性能损耗极低 spark_df = spark.createDataFrame(pa_table) spark_df.show()
2. 读取指定路径下所有Parquet文件
# 开启Spark Arrow支持 spark.conf.set("spark.sql.execution.arrow.enabled", "true") import io import pyarrow.parquet as pq import pyarrow as pa res = autorefresh_session.resource('s3') bucket = res.Bucket(name='mybucket') prefix_path = 'my/path/' pa_table_list = [] # 遍历前缀下所有对象 for obj in bucket.objects.filter(Prefix=prefix_path): # 过滤有效parquet文件,跳过目录占位符 if obj.key.endswith('.parquet') and not obj.key.endswith('/'): body = io.BytesIO(obj.get()['Body'].read()) pa_table_list.append(pq.read_table(body)) # 合并所有Arrow表为全量表 full_pa_table = pa.concat_tables(pa_table_list) # 转换为Spark DataFrame spark_df = spark.createDataFrame(full_pa_table) spark_df.show()
内容的提问来源于stack exchange,提问作者pkp
相关产品推荐
相关产品推荐

