Fugue项目快速入门:10分钟掌握核心API
2025-06-10 21:30:45作者:滑思眉Philip
项目概述
Fugue是一个旨在简化大数据处理流程的开源项目,它通过提供统一的接口让用户能够轻松地在不同计算引擎(如Spark、Dask、Ray等)上执行分布式计算。本文将带您快速了解Fugue的核心API功能,帮助数据从业者快速上手使用。
适用人群
Fugue特别适合以下三类用户:
- 需要将Python或Pandas编写的业务逻辑扩展到更大数据集的数据科学家
- 希望通过分布式计算并行化现有代码的数据从业者
- 希望减少Spark/Dask/Ray代码维护和测试工作量的数据团队
环境准备
首先我们需要初始化一个Spark会话,后续示例会用到:
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
核心功能:transform()函数
Fugue最核心的功能是transform()函数,它能让用户轻松地将Pandas或Python代码扩展到分布式执行环境,而只需做最小的代码修改。
基础示例:模型预测
让我们通过一个机器学习预测的例子来演示:
- 首先训练一个简单的线性回归模型:
import pandas as pd
import numpy as np
from sklearn.linear_model import LinearRegression
X = pd.DataFrame({"x_1": [1, 1, 2, 2], "x_2":[1, 2, 2, 3]})
y = np.dot(X, np.array([1, 2])) + 3
reg = LinearRegression().fit(X, y)
- 然后定义一个预测函数:
def predict(df: pd.DataFrame, model: LinearRegression) -> pd.DataFrame:
"""使用预训练模型进行预测"""
return df.assign(predicted=model.predict(df))
# 测试数据
input_df = pd.DataFrame({"x_1": [3, 4, 6, 6], "x_2":[3, 3, 6, 6]})
# 本地测试
predict(input_df, reg)
- 现在使用Fugue将这个函数扩展到Spark执行:
from fugue import transform
result = transform(
df=input_df,
using=predict,
schema="*,predicted:double",
params=dict(model=reg),
engine=spark
)
print(type(result)) # 输出: <class 'pyspark.sql.dataframe.DataFrame'>
result.show()
transform()函数参数解析
df: 输入DataFrame(可以是Pandas、Spark、Dask或Ray DataFrame)using: 要应用的Python函数schema: 输出结果的schema定义params: 传递给函数的参数字典engine: 执行引擎(Pandas、Spark、Dask或Ray)
执行引擎选择
Fugue支持多种执行引擎,使用方式非常灵活:
# 使用Spark
transform(df, fn, ..., engine=spark_session) # 输出Spark DataFrame
# 使用Dask
transform(df, fn, ..., engine=dask_client) # 输出Dask DataFrame
# 使用Ray
transform(df, fn, ..., engine="ray") # 输出Ray Dataset
如果不指定engine参数,Fugue会根据输入DataFrame的类型自动选择执行引擎:
transform(df, fn, ...) # 使用Pandas
transform(spark_df, fn, ...) # 使用Spark
transform(dask_df, fn, ...) # 使用Dask
transform(ray_df, fn, ...) # 使用Ray
本地结果返回
默认情况下,Fugue不会将分布式DataFrame转换为本地Pandas DataFrame。如果需要本地结果,可以设置as_local=True:
local_result = transform(
df=input_df,
using=predict,
schema="*,predicted:double",
params=dict(model=reg),
engine=spark,
as_local=True
)
print(type(local_result)) # 输出: <class 'pandas.core.frame.DataFrame'>
注意:对于大数据集,不建议返回本地DataFrame,可能会造成驱动程序内存不足。
类型提示与转换
Fugue通过函数类型提示来指导数据转换。前面的例子使用了pd.DataFrame作为输入输出类型,但Fugue也支持其他格式:
- 使用字典列表作为输入输出:
from typing import List, Dict, Any
def add_row2(df: List[Dict[str,Any]]) -> List[Dict[str,Any]]:
result = []
for row in df:
row["total"] = row["a"] + row["b"] + row["c"]
if row["total"] < 10:
result.append(row)
return result
- 使用列表的列表作为输入输出:
from typing import List, Iterable, Any
def add_row3(df: List[List[Any]]) -> Iterable[List[Any]]:
for row in df:
row.append(sum(row))
if row[-1] < 10:
yield row
这些函数都可以直接使用transform()函数在分布式环境中执行,Fugue会自动处理类型转换。
总结
通过本文,我们快速了解了Fugue项目的核心功能:
- 使用
transform()函数轻松将Pandas/Python代码扩展到分布式环境 - 支持多种执行引擎(Spark、Dask、Ray)的无缝切换
- 灵活的类型系统支持多种数据格式
- 简化分布式代码的测试和维护
Fugue的强大之处在于它让开发者可以专注于业务逻辑,而不必担心底层分布式计算的复杂性。对于需要处理大数据的Python开发者来说,Fugue是一个非常值得尝试的工具。
登录后查看全文
热门项目推荐
GLM-5智谱 AI 正式发布 GLM-5,旨在应对复杂系统工程和长时域智能体任务。Jinja00
GLM-5-w4a8GLM-5-w4a8基于混合专家架构,专为复杂系统工程与长周期智能体任务设计。支持单/多节点部署,适配Atlas 800T A3,采用w4a8量化技术,结合vLLM推理优化,高效平衡性能与精度,助力智能应用开发Jinja00
请把这个活动推给顶尖程序员😎本次活动专为懂行的顶尖程序员量身打造,聚焦AtomGit首发开源模型的实际应用与深度测评,拒绝大众化浅层体验,邀请具备扎实技术功底、开源经验或模型测评能力的顶尖开发者,深度参与模型体验、性能测评,通过发布技术帖子、提交测评报告、上传实践项目成果等形式,挖掘模型核心价值,共建AtomGit开源模型生态,彰显顶尖程序员的技术洞察力与实践能力。00
Kimi-K2.5Kimi K2.5 是一款开源的原生多模态智能体模型,它在 Kimi-K2-Base 的基础上,通过对约 15 万亿混合视觉和文本 tokens 进行持续预训练构建而成。该模型将视觉与语言理解、高级智能体能力、即时模式与思考模式,以及对话式与智能体范式无缝融合。Python00
MiniMax-M2.5MiniMax-M2.5开源模型,经数十万复杂环境强化训练,在代码生成、工具调用、办公自动化等经济价值任务中表现卓越。SWE-Bench Verified得分80.2%,Multi-SWE-Bench达51.3%,BrowseComp获76.3%。推理速度比M2.1快37%,与Claude Opus 4.6相当,每小时仅需0.3-1美元,成本仅为同类模型1/10-1/20,为智能应用开发提供高效经济选择。【此简介由AI生成】Python00
Qwen3.5Qwen3.5 昇腾 vLLM 部署教程。Qwen3.5 是 Qwen 系列最新的旗舰多模态模型,采用 MoE(混合专家)架构,在保持强大模型能力的同时显著降低了推理成本。00- RRing-2.5-1TRing-2.5-1T:全球首个基于混合线性注意力架构的开源万亿参数思考模型。Python00
热门内容推荐
最新内容推荐
Degrees of Lewdity中文汉化终极指南:零基础玩家必看的完整教程Unity游戏翻译神器:XUnity Auto Translator 完整使用指南PythonWin7终极指南:在Windows 7上轻松安装Python 3.9+终极macOS键盘定制指南:用Karabiner-Elements提升10倍效率Pandas数据分析实战指南:从零基础到数据处理高手 Qwen3-235B-FP8震撼升级:256K上下文+22B激活参数7步搞定机械键盘PCB设计:从零开始打造你的专属键盘终极WeMod专业版解锁指南:3步免费获取完整高级功能DeepSeek-R1-Distill-Qwen-32B技术揭秘:小模型如何实现大模型性能突破音频修复终极指南:让每一段受损声音重获新生
项目优选
收起
deepin linux kernel
C
27
11
OpenHarmony documentation | OpenHarmony开发者文档
Dockerfile
569
3.84 K
Ascend Extension for PyTorch
Python
379
454
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
893
677
暂无简介
Dart
802
199
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
350
205
昇腾LLM分布式训练框架
Python
118
147
Nop Platform 2.0是基于可逆计算理论实现的采用面向语言编程范式的新一代低代码开发平台,包含基于全新原理从零开始研发的GraphQL引擎、ORM引擎、工作流引擎、报表引擎、规则引擎、批处理引引擎等完整设计。nop-entropy是它的后端部分,采用java语言实现,可选择集成Spring框架或者Quarkus框架。中小企业可以免费商用
Java
12
1
🔥LeetCode solutions in any programming language | 多种编程语言实现 LeetCode、《剑指 Offer(第 2 版)》、《程序员面试金典(第 6 版)》题解
Java
68
20
🎉 (RuoYi)官方仓库 基于SpringBoot,Spring Security,JWT,Vue3 & Vite、Element Plus 的前后端分离权限管理系统
Vue
1.37 K
781