You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 08:36:05