如何使用Ray Clusters替代Dask并行读取Parquet并完成合并计算?
Ray集群下并行读取Parquet并合并的代码修改
问题背景
原本使用Dask读取两个Parquet文件并合并,现在想改用Ray集群提升性能,但遇到代码错误:
- 原Dask代码:
def read_and_merge_parquets(): df1 = dd.read_parquet(path='parquet1.parquet').compute() df2 = dd.read_parquet(path='parquet2.parquet').compute() merged_df = df2.merge(df1, on="id", how="left") print(merged_df)
- 尝试的Ray代码及错误:
ray.init(address='auto') def read_and_merge_parquets(): df1 = ray.data.read_parquet(paths='parquet1.parquet').compute() df2 = ray.data.read_parquet(paths='parquet2.parquet').compute() merged_df = df2.merge(df1, on="id", how="left") print(merged_df)
报错:AttributeError: 'Dataset' object has no attribute 'compute',同时尝试enable_dask_on_ray()时脚本仅在头节点运行,无法利用3个工作节点。
解决方案
错误原因说明
Ray Data的ray.data.read_parquet返回的是分布式Dataset对象,它没有compute()方法——这是Dask专属的API。如果要转换成本地Pandas DataFrame,需要用to_pandas();但如果想利用Ray集群的并行能力,应该直接使用Ray Data的分布式操作,避免把全量数据拉回本地。
方案1:利用Ray Data分布式Join(推荐,充分利用集群)
这种方式全程在Ray集群上并行执行读取、合并,不会把全量数据拉到头节点,能最大化利用3个工作节点的资源:
import ray # 连接到Ray集群,address='auto'会自动发现集群地址 ray.init(address='auto') def read_and_merge_parquets(): # 并行读取Parquet文件,返回分布式Dataset ds1 = ray.data.read_parquet('parquet1.parquet') ds2 = ray.data.read_parquet('parquet2.parquet') # 执行分布式左连接,on指定关联键,how指定连接类型 merged_ds = ds2.join(ds1, on="id", join_type="left") # 如果需要查看结果,可以选择: # 1. 取前N条数据打印(避免全量拉取) print(merged_ds.take(10)) # 2. 转换成Pandas DataFrame(仅适合数据量不大的情况) # merged_df = merged_ds.to_pandas() # print(merged_df) if __name__ == "__main__": read_and_merge_parquets()
方案2:转换为Pandas DataFrame(适合小数据场景)
如果数据量较小,一定要转换成本地Pandas对象处理,把compute()替换为to_pandas()即可:
import ray ray.init(address='auto') def read_and_merge_parquets(): df1 = ray.data.read_parquet('parquet1.parquet').to_pandas() df2 = ray.data.read_parquet('parquet2.parquet').to_pandas() merged_df = df2.merge(df1, on="id", how="left") print(merged_df) if __name__ == "__main__": read_and_merge_parquets()
关键注意事项
- 必须把执行逻辑放在
if __name__ == "__main__"块中,这是Ray分布式代码的标准要求,避免工作节点重复初始化代码。 - 方案1的
join操作是分布式执行的,数据会在集群节点间分片处理,不会集中到头节点,能真正利用3个工作节点的算力。 - 如果Parquet文件是分布式存储(如S3、HDFS),直接传入路径即可,Ray Data会自动并行读取分片。
内容的提问来源于stack exchange,提问作者Akham01
相关产品推荐
相关产品推荐

