首页
/ Fugue项目教程:深入理解执行引擎(Execution Engine)的使用

Fugue项目教程:深入理解执行引擎(Execution Engine)的使用

2025-06-10 23:06:33作者:舒璇辛Bertina

引言

在分布式计算领域,执行引擎(Execution Engine)是数据处理的核心组件。Fugue作为一个统一的分布式计算框架,提供了多种执行引擎的支持,使开发者能够轻松地在不同计算环境间切换。本文将详细介绍Fugue中执行引擎的使用方法,帮助开发者掌握这一关键功能。

执行引擎概述

Fugue支持多种执行引擎,包括Spark、Dask和Ray等。通过统一的API接口,开发者可以无缝切换不同的计算后端,而无需重写业务逻辑代码。这种设计极大地提高了代码的可移植性和开发效率。

基础设置

在开始之前,我们先准备一个简单的示例环境:

import pandas as pd
from fugue import transform 

# 创建示例数据
df = pd.DataFrame({"col1": [1,2,3,4], "col2": [1,2,3,4]})

# 定义转换函数,添加新列col3
# schema: *, col3:int
def add_cols(df:pd.DataFrame) -> pd.DataFrame:
    return df.assign(col3 = df['col1'] + df['col2'])

这个示例展示了如何创建一个简单的Pandas DataFrame,并定义一个函数来添加新列。注意函数定义中的schema提示,这是Fugue的一个有用特性,可以明确指定输出数据的结构。

通过字符串指定执行引擎

最简单的引擎指定方式是直接传递引擎名称字符串。Fugue会自动在本地启动相应的计算引擎,并利用所有可用的计算资源。

Spark引擎示例

spark_df = transform(df, add_cols, engine="spark")
spark_df.show()

执行结果将显示包含新列col3的DataFrame。Spark引擎适合处理大规模数据集,提供了强大的分布式计算能力。

Dask引擎示例

dask_df = transform(df, add_cols, engine="dask")
dask_df.compute().head()

Dask引擎提供了类似Pandas的API,但支持并行计算,特别适合中等规模数据的处理。

Ray引擎示例

ray_df = transform(df, add_cols, engine="ray")
ray_df.show(5)

Ray是一个新兴的分布式计算框架,特别适合机器学习和AI工作负载,提供了低延迟和高吞吐量的计算能力。

通过Client或Session对象指定引擎

除了字符串方式,Fugue还支持直接传递已创建的引擎会话对象,这种方式提供了更大的灵活性。

Spark会话示例

from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()

spark_df = transform(df, add_cols, engine=spark)
spark_df.show()

这种方式允许重用现有的Spark会话,避免重复创建的开销。

Dask Client示例

from distributed import Client
dask_client = Client()

dask_df = transform(df, add_cols, engine=dask_client)
dask_df.compute().head()

通过Dask Client对象,可以更精细地控制计算资源分配。

Ray初始化

import ray
ray.init(ignore_reinit_error=True)

ray_df = transform(df, add_cols, engine="ray")
ray_df.show(5)

Ray的初始化方式略有不同,但同样简单直观。

集群连接与配置

Fugue还支持直接连接到远程计算集群,这为大规模分布式计算提供了便利。

集群连接示例

# Databricks集群
transform(df, add_cols, engine="db", engine_conf=conf)

# Coiled集群
transform(df, add_cols, engine="coiled:my_cluster")

# Anyscale集群
transform(df, add_cols, engine="anyscale://project/cluster-1")

这些连接方式需要额外的认证配置,但提供了强大的云端计算能力。

引擎配置选项

Fugue允许通过engine_conf参数对执行引擎进行详细配置。

配置示例

spark_df = transform(df, 
                     add_cols, 
                     engine=spark, 
                     engine_conf={"fugue.spark.use_pandas_udf":True})
spark_df.show(2)

这个示例启用了Spark的Pandas UDF功能,可以优化某些类型的数据处理性能。

最佳实践与建议

  1. 开发阶段:建议先使用本地模式(如Dask或Ray本地模式)进行开发和测试,确认逻辑正确后再切换到分布式环境。

  2. 生产环境:根据数据规模和计算需求选择合适的引擎。大规模数据处理推荐Spark,机器学习任务可考虑Ray。

  3. 性能调优:合理设置分区数和执行器资源,可以显著提高计算效率。

  4. 错误处理:熟悉不同引擎的日志格式,便于快速定位问题。

总结

Fugue的执行引擎抽象层为分布式计算提供了极大的便利。通过本文介绍的各种引擎指定方式,开发者可以根据项目需求灵活选择最适合的计算后端。无论是本地开发还是云端部署,Fugue都能提供一致且高效的开发体验。

掌握这些执行引擎的使用方法,将帮助您更好地利用Fugue的强大功能,构建高效可靠的数据处理流程。

登录后查看全文
热门项目推荐

项目优选

收起
openHiTLS-examplesopenHiTLS-examples
本仓将为广大高校开发者提供开源实践和创新开发平台,收集和展示openHiTLS示例代码及创新应用,欢迎大家投稿,让全世界看到您的精巧密码实现设计,也让更多人通过您的优秀成果,理解、喜爱上密码技术。
C
54
468
kernelkernel
deepin linux kernel
C
22
5
nop-entropynop-entropy
Nop Platform 2.0是基于可逆计算理论实现的采用面向语言编程范式的新一代低代码开发平台,包含基于全新原理从零开始研发的GraphQL引擎、ORM引擎、工作流引擎、报表引擎、规则引擎、批处理引引擎等完整设计。nop-entropy是它的后端部分,采用java语言实现,可选择集成Spring框架或者Quarkus框架。中小企业可以免费商用
Java
7
0
RuoYi-Vue3RuoYi-Vue3
🎉 (RuoYi)官方仓库 基于SpringBoot,Spring Security,JWT,Vue3 & Vite、Element Plus 的前后端分离权限管理系统
Vue
879
517
Cangjie-ExamplesCangjie-Examples
本仓将收集和展示高质量的仓颉示例代码,欢迎大家投稿,让全世界看到您的妙趣设计,也让更多人通过您的编码理解和喜爱仓颉语言。
Cangjie
336
1.1 K
ohos_react_nativeohos_react_native
React Native鸿蒙化仓库
C++
180
264
cjoycjoy
一个高性能、可扩展、轻量、省心的仓颉Web框架。Rest, 宏路由,Json, 中间件,参数绑定与校验,文件上传下载,MCP......
Cangjie
87
14
CangjieCommunityCangjieCommunity
为仓颉编程语言开发者打造活跃、开放、高质量的社区环境
Markdown
1.08 K
0
openHiTLSopenHiTLS
旨在打造算法先进、性能卓越、高效敏捷、安全可靠的密码套件,通过轻量级、可剪裁的软件技术架构满足各行业不同场景的多样化要求,让密码技术应用更简单,同时探索后量子等先进算法创新实践,构建密码前沿技术底座!
C
359
381
cherry-studiocherry-studio
🍒 Cherry Studio 是一款支持多个 LLM 提供商的桌面客户端
TypeScript
612
60