当前位置: 首页 > news >正文

【daft框架】和ray分布式计算的结合运行自定义函数

daft的框架

主要分成python和raft两部分

daft在ray上如何运行udf

采用分布式执行框架

  • Client 端: RemoteFlotillaRunner 负责把物理计划切成任务,分发到各个节点
  • Worker 端: 每个 Ray 节点上只跑一个 RaySwordfishActor
  • 内部调度: Actor 内部有任务队列,根据 UDF 声明的资源 (num_gpus/num_cpus) 调度任务

具体流程

步骤1 ray.init()

  • 连接 Ray 集群,这一步就是 Ray 本身的逻辑,和 Daft 无关
  • Daft 复用你已经初始化好的 Ray,不需要自己重新初始化

步骤2 daft.set_runner_ray()

这一步才是 Daft 启动 Actor 的地方:

def set_runner_ray():
# → 创建 RemoteFlotillaRunner
# → 在每个 Ray 节点上启动一个 RaySwordfishActor
# → 这些 Actors 启动好之后,就一直运行,等待任务

# Daft 在 set_runner_ray 的时候,在每个 Ray 节点启动一个 Actor:@ray.remoteclassRaySwordfishActor:def__init__(self):# 这里启动好,一直活着self.task_queue=...self.resource_scheduler=...self.instantiated_udfs={}# 缓存已经实例化的 UDF

关键点: 此时只启动了一个空的 RaySwordfishActor 每个节点,Actor 只是个空壳,里面还没有任何 UDF。

步骤3 UDF 定义与注册阶段

@daft.func(num_gpus=0.5,concurrency=2)defmy_udf(image):returnmodel.predict(image)df=df.with_column("prediction",my_udf(col("image")))

在 Python 层定义 UDF

发生了什么:

  • @daft.func/@daft.cls 将你的函数/类包装成 Daft 内部的 UDF 对象
  • 资源配置(num_gpus/num_cpus/ray_options)被存在 UDF 对象里
  • UDF 信息被注册到 Daft 的函数注册表 。不申请资源,不启动任何东西,只是保存信息

步骤4 查询规划阶段

当执行 df.collect()

result=df.collect()

发生了什么:

  1. Daft 从逻辑计划 → 优化 → 生成物理执行计划

  2. 物理计划会把计算切分成多个任务块,每个任务块处理一批数据

  3. Flotilla 调度器知道哪些 UDF 需要什么资源

  4. 任务提交与调度

物理计划生成后,Flotilla 把任务发给各个节点的 RaySwordfishActor:

Client → RemoteFlotillaRunner → 分推任务 → 各个 RaySwordfishActor

调度逻辑:

  • Daft 会根据每个 UDF 声明的资源需求(num_gpus/num_cpus)做内部调度
  • concurrency=N 决定同一个 Actor 里最多同时跑几个该 UDF 的任务
  • 对于 GPU: 如果你声明 num_gpus=0.5,同一个 Actor 可以并行跑 2 个,共享同一块 GPU,这是 Daft 比原生 Ray 好的地方

步骤5. UDF 执行

当 RaySwordfishActor 收到一个 UDF 任务:

  1. 反序列化: 从任务描述中拿到 UDF 和输入数据
  2. 实例化: 如果是类 UDF(@daft.cls),实例化你的类(只实例化一次,复用实例)
  3. 执行: 调用你的 UDF 处理输入 batch
  4. 序列化: 把输出结果序列化,传回给下游或者 client

关键优化:

  • UDF 实例复用: 相同 UDF 只实例化一次,不会每个任务都新建,节省初始化开销(比如模型只加载一次到 GPU)
  • 批处理: Daft 会把数据攒成批再给你的 UDF,提升利用率
  • 内存管理: 大批次会自动拆分,避免 OOM

核心代码位置

  • Flotilla 入口: daft/runners/flotilla.py → RemoteFlotillaRunner
  • Swordfish Actor: daft/runners/swordfish/actor.py → RaySwordfishActor
  • UDF v2 实现: daft/udf/udf_v2.py
  • Shuffle 实现: src/daft-shuffles/ (Rust)
http://www.cnnetsun.cn/news/1877332.html

相关文章:

  • 【青少年CTF S1·2026 公益赛】时间胶囊留言板
  • 属性图:节点、边与属性的图模型
  • 【CTF | pwn篇】从ctfshow入门到进阶:栈溢出实战技巧全解析
  • TexLive极简安装法:5分钟搞定基础版+中英文支持(附磁盘空间不足解决方案)
  • 从理论到硅片:二值化CNN在FPGA上的高效部署实践
  • Spring AI Alibaba 1.1
  • c++读取整个文件到字符串方法 c++如何一次性读取文件
  • ESOP软硬件方案为3C产线提供全方位防错保障
  • AEB标配率接近8成!强规临近,存量红利在哪
  • devops系列(一) Nginx 反向代理与负载均衡:一台服务器扛不住怎么办
  • MCP与Agent协同的智能体架构设计
  • 代码随想录算法训练营第七天补卡| leetcode144\145\94 二叉树的左右中序遍历
  • ViPER4Windows音频补丁工具:Windows 10/11系统兼容性问题终极解决方案
  • 深度解析ImageNet分类任务中的卷积神经网络架构优化策略
  • 【数据结构与算法】第43篇:Trie树(前缀树/字典树)
  • 通义千问3-4B真实体验:本地部署生成测试用例,效率提升实测
  • 压力测试下的心态管理:当线上告警电话响起时
  • 别再瞎测了!手把手教你用泰克/安捷伦示波器搞定USB2.0信号质量(Device/Hub模式实战)
  • 从HTTP到K8s探针:彻底搞懂HTTPGet健康检查
  • 全网最全:零基础学深度学习需要学哪些框架?PyTorch 和 TensorFlow 选哪个?
  • 告别手写脚本!用Frida-Trace自动Hook Android App的Java方法(附实战Demo)
  • ToClaw真的能让AI Agent落地吗?先看清它的代价与边界
  • Python语言的12个基础知识点小结
  • 基于STM32的智能家居安防系统设计与实现
  • LangChain4j 1.0.0-beta2踩坑记:从社区版DashScope依赖到SpringBoot自动配置的完整避坑指南
  • 扣子(Coze)实战:10万+治愈奶奶图文,Coze一键生成
  • Simulink信号解析避坑指南:为什么你的‘蓝色鱼叉’图标不出现?
  • [Unity] ShaderGraph实战:动态水面倒影与镜面反射效果优化
  • SDXL 1.0电影级绘图工坊:Mathtype公式渲染与科学图表生成
  • SQL如何获取分组最后一条数据_LAST_VALUE的滑动窗口陷阱