Fugue快速入门:10分钟学会用Python代码跨Spark、Dask、Ray运行
【免费下载链接】fugue A unified interface for distributed computing. Fugue executes SQL, Python, Pandas, and Polars code on Spark, Dask and Ray without any rewrites. 项目地址: https://gitcode.com/gh_mirrors/fu/fugue
Fugue是一个统一的分布式计算接口,让你无需重写代码就能在Spark、Dask和Ray等分布式框架上执行SQL、Python、Pandas和Polars代码。无论是数据处理新手还是有经验的开发者,都能通过Fugue轻松实现代码的跨平台运行,极大提升工作效率。
🚀 为什么选择Fugue?
在大数据处理领域,不同的分布式框架(如Spark、Dask、Ray)各有优势,但它们的API差异往往让开发者头疼。Fugue通过提供统一的抽象层,完美解决了这一痛点:
- 代码复用:一次编写,多框架运行,无需为每个框架单独适配
- 学习成本低:使用熟悉的Python/Pandas语法,无需深入学习各框架API
- 灵活性高:随时切换执行引擎,根据任务需求选择最优框架
Fugue架构图:展示了Fugue如何通过抽象层连接不同的计算引擎
🔧 安装Fugue
安装Fugue非常简单,通过pip命令即可完成。根据你的需求,可以选择安装不同的扩展组件:
# 基础安装
pip install fugue
# 如需支持Spark
pip install fugue[spark]
# 如需支持Dask
pip install fugue[dask]
# 如需支持Ray
pip install fugue[ray]
# 完整安装(包含所有计算引擎)
pip install fugue[all]
详细依赖信息可查看项目根目录下的requirements.txt文件
🔍 Fugue核心概念
Fugue的核心设计围绕"抽象"和"扩展"两个关键词展开:
抽象层设计
Fugue通过ExecutionEngine抽象层,将用户代码与底层计算引擎解耦。用户只需关注业务逻辑,无需关心具体执行细节。
扩展机制
Fugue提供了丰富的扩展接口,让你可以轻松扩展功能:
Fugue扩展架构:展示了Creator、Processor、Outputter等核心扩展组件
主要扩展组件包括:
- Creator:数据创建组件
- Processor:数据处理组件
- Outputter:数据输出组件
- Transformer:数据转换组件
这些扩展可以在extensions/目录中找到实现源码。
📝 快速上手示例
下面通过一个简单示例,展示如何使用Fugue在不同引擎上运行相同的代码。
1. 基础Pandas代码
首先,我们编写一个简单的Pandas数据处理函数:
import pandas as pd
def process_data(df: pd.DataFrame) -> pd.DataFrame:
return df[df["value"] > 0].groupby("category").mean()
2. 使用Fugue包装
只需添加Fugue的装饰器,即可将普通函数转换为分布式函数:
from fugue import transformer
@transformer
def process_data(df: pd.DataFrame) -> pd.DataFrame:
return df[df["value"] > 0].groupby("category").mean()
3. 在不同引擎上运行
from fugue import FugueWorkflow
from fugue_spark import SparkExecutionEngine
from fugue_dask import DaskExecutionEngine
from fugue_ray import RayExecutionEngine
# 创建测试数据
data = [("A", 1), ("A", -2), ("B", 3), ("B", 4)]
df = pd.DataFrame(data, columns=["category", "value"])
# 默认引擎(本地Pandas)
with FugueWorkflow() as wf:
result = wf.df(df).transform(process_data)
result.show()
# Spark引擎
with FugueWorkflow(SparkExecutionEngine) as wf:
result = wf.df(df).transform(process_data)
result.show()
# Dask引擎
with FugueWorkflow(DaskExecutionEngine) as wf:
result = wf.df(df).transform(process_data)
result.show()
# Ray引擎
with FugueWorkflow(RayExecutionEngine) as wf:
result = wf.df(df).transform(process_data)
result.show()
以上代码无需修改业务逻辑,即可在不同的分布式引擎上运行,真正实现了"一次编写,到处运行"。
📚 学习资源
- 官方文档:项目中的docs/目录包含完整的使用文档和API参考
- 示例代码:tests/目录中有大量测试用例,可作为学习参考
- 扩展模块:Fugue提供了多个扩展模块,如fugue_spark/、fugue_dask/和fugue_ray/,分别对应不同的计算引擎
💡 总结
Fugue为Python数据处理提供了一个强大而灵活的统一接口,让你能够轻松应对不同的分布式计算场景。通过Fugue,你可以:
无论你是数据科学家、分析师还是工程师,Fugue都能帮助你更高效地处理数据,让你的Python代码发挥更大价值!
如果你想深入了解Fugue的更多功能,可以通过以下命令克隆项目源码进行探索:
git clone https://gitcode.com/gh_mirrors/fu/fugue
【免费下载链接】fugue A unified interface for distributed computing. Fugue executes SQL, Python, Pandas, and Polars code on Spark, Dask and Ray without any rewrites. 项目地址: https://gitcode.com/gh_mirrors/fu/fugue
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考



