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
相关产品推荐
相关产品推荐

