从S3指定日期文件夹读取过滤Parquet文件的替代方案(R/Python/Databricks)
解决Parquet大文件内存不足问题的替代实现方案
Python 实现方式
pandas + fastparquet 分块+谓词下推
借助fastparquet引擎支持的谓词下推,提前过滤目标ID,避免全量加载数据,同时通过分块读取控制内存占用:import pandas as pd # 配置好AWS凭证后指定目标日期的S3路径 s3_path = "s3://your-bucket/date=2024-01-01/" target_ids = [1001, 1002, 1003] # 分块读取并过滤 chunk_list = [] for chunk in pd.read_parquet( s3_path, engine='fastparquet', filters=[('id', 'in', target_ids)], chunksize=100000 ): chunk_list.append(chunk) # 合并为最终DataFrame final_df = pd.concat(chunk_list, ignore_index=True)Dask 分布式分块处理
Dask会自动将数据拆分为多个块并行处理,无需一次性加载全量数据到内存:import dask.dataframe as dd s3_path = "s3://your-bucket/date=2024-01-01/" target_ids = [1001, 1002, 1003] # 读取目标文件夹下所有Parquet文件 ddf = dd.read_parquet(s3_path) # 过滤特定ID filtered_ddf = ddf[ddf['id'].isin(target_ids)] # 按需转为pandas DataFrame,也可直接用Dask进行后续分析 final_df = filtered_ddf.compute()PySpark 分布式读取与过滤
Spark的分布式架构天然适配大文件处理,谓词下推会在存储层执行过滤,大幅减少数据传输量:from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ParquetIDFilter").getOrCreate() s3_path = "s3://your-bucket/date=2024-01-01/" target_ids = [1001, 1002, 1003] # 读取Parquet文件 df = spark.read.parquet(s3_path) # 过滤目标ID filtered_df = df.filter(df.id.isin(target_ids)) # 按需转为pandas DataFrame final_df = filtered_df.toPandas()
R 实现方式
sparklyr 对接Spark处理
通过sparklyr连接Spark集群,利用分布式能力规避本地内存限制:library(sparklyr) # 初始化Spark连接(生产环境替换为集群地址) sc <- spark_connect(master = "local") s3_path <- "s3://your-bucket/date=2024-01-01/" target_ids <- c(1001, 1002, 1003) # 读取目标日期的Parquet文件夹 df <- spark_read_parquet(sc, path = s3_path) # 过滤特定ID filtered_df <- df %>% filter(id %in% target_ids) # 按需提取到本地DataFrame final_df <- collect(filtered_df) # 关闭Spark连接 spark_disconnect(sc)duckdb + dplyr 谓词下推查询
将Parquet文件视为数据库表,通过duckdb执行谓词下推,仅加载需要的数据:library(duckdb) library(dplyr) s3_path <- "s3://your-bucket/date=2024-01-01/" target_ids <- c(1001, 1002, 1003) # 连接duckdb con <- dbConnect(duckdb()) # 注册Parquet文件夹为视图 dbExecute(con, paste0("CREATE VIEW parquet_data AS SELECT * FROM parquet_scan('", s3_path, "')")) # 过滤查询目标ID final_df <- dbGetQuery( con, paste0("SELECT * FROM parquet_data WHERE id IN (", paste(target_ids, collapse = ","), ")") ) # 关闭连接 dbDisconnect(con)
Databricks 实现方式
Spark SQL 直接查询
利用Databricks的分布式Spark集群,直接通过SQL读取并过滤数据:SELECT * FROM parquet.`s3://your-bucket/date=2024-01-01/` WHERE id IN (1001, 1002, 1003)执行后可直接将结果转为DataFrame,或保存至其他存储系统。
PySpark 代码处理
在Databricks notebook中通过PySpark代码完成读取与过滤:s3_path = "s3://your-bucket/date=2024-01-01/" target_ids = [1001, 1002, 1003] # 读取目标日期的Parquet文件 df = spark.read.parquet(s3_path) # 过滤特定ID filtered_df = df.filter(df.id.isin(target_ids)) # 可直接使用filtered_df进行后续分析,或按需转为Pandas DataFrame final_df = filtered_df.toPandas()
内容的提问来源于stack exchange,提问作者Datamaniac
相关产品推荐
相关产品推荐

