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

Ray是否提供声明式/函数式接口实现远程函数的迭代器映射?

当前代码

#!/usr/bin/env python3
# encoding: utf-8
"""Demonstration of Ray parallelism"""
import ray
from typing import Iterator
ray.init()
@ray.remote
def square(n:int)->int:
    return n*n

references: Iterator[ray.ObjectRef] = map(lambda val: square.remote(val),
                                                    range(10))
ray.get([*references])
ray.shutdown()

这段代码本质上实现了Ray版本的map(square, range(10))功能。

问题

对于这种标准且常见的模式,上述实现过于冗长。请问Ray是否提供更具声明式/函数式风格的API来实现上述功能?除map外,最好还支持filter、reduce等操作。


Ray确实提供了更简洁的声明式/函数式API来处理这类并行操作,核心是Ray Data模块(原Ray Dataset),它专为批量数据处理设计,支持map、filter、reduce等常见函数式操作,写法贴近原生Python风格,同时自动完成并行化执行。

并行map实现示例

#!/usr/bin/env python3
import ray

ray.init()

def square(n: int) -> int:
    return n * n

# 创建数据集并执行并行map
ds = ray.data.from_items(range(10))
result = ds.map(square).take_all()
print(result)

ray.shutdown()

filter操作示例

筛选偶数后再执行平方:

import ray

ray.init()

def square(n: int) -> int:
    return n * n

def is_even(n: int) -> bool:
    return n % 2 == 0

ds = ray.data.from_items(range(10))
result = ds.filter(is_even).map(square).take_all()
print(result)  # 输出: [0, 4, 16, 36, 64]

ray.shutdown()

reduce操作示例

计算所有平方值的总和:

import ray

ray.init()

def square(n: int) -> int:
    return n * n

def sum_accumulator(a: int, b: int) -> int:
    return a + b

ds = ray.data.from_items(range(10))
total = ds.map(square).reduce(sum_accumulator)
print(total)  # 输出: 285

ray.shutdown()

额外优势

  • 无需手动使用@ray.remote装饰器和remote()调用,代码更简洁
  • 自动处理任务调度、并行执行,无需手动管理ObjectRef和ray.get()
  • 支持大规模分布式数据处理,自动完成数据分区和资源调度

内容的提问来源于stack exchange,提问作者Della

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 21:08:12