Daft on Ray 实战指南:从本地 Ray 集群部署到自动扩缩容(Flotilla 分布式执行深度解析)

发布时间:2026/9/17 20:10:40
Daft on Ray 实战指南:从本地 Ray 集群部署到自动扩缩容(Flotilla 分布式执行深度解析) Daft on Ray 实战指南从本地 Ray 集群部署到自动扩缩容Flotilla 分布式执行深度解析【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft本文基于 Daft 官方文档 Ray 运行指南 整理并深入扩展系统讲解 Daft 如何借助 Ray 分布式框架执行 DataFrame 查询包括本地单机 Ray 集群的快速搭建、远程集群连接、Ray Client 与 Ray Jobs 两种接入方式以及 Daft 内置的 autoscaler 扩缩容scale-up / scale-in机制。读完本文你将能够把 Daft 部署到各类 Ray 环境中理解daft.set_runner_ray()的完整参数语义并从源码层面掌握 Daft 如何与 Ray autoscaler 协商资源。Daft 在 Ray 上的执行架构Flotilla 与 Swordfish Worker在动手配置之前先理解 Daft 与 Ray 的集成方式。从源码结构看Daft 的 Ray 后端由三层组成Python 侧入口daft/runners/init.py 中的set_runner_ray()负责把用户的连接参数地址、扩缩容策略等配置进执行上下文并将部分配置写入环境变量传递给 Rust 调度器Flotilla 运行时daft/runners/flotilla.py 中的FlotillaRunner/RemoteFlotillaRunner是一个 Ray actor负责在集群中拉起和跟踪 worker并暴露start_ray_workers、try_autoscale、clear_autoscaling_requests、get_head_node_id等关键函数Rust 侧调度核心src/daft-distributed/src/python/ray/worker_manager.rs 中的RayWorkerManager实现了 worker 发现、任务提交、自动扩缩容autoscale与空闲 worker 回收retire的全部决策逻辑。每个 Ray 节点上运行一个 Swordfish worker即 daft/runners/flotilla.py 中start_ray_workers按节点创建的 Ray actor使用NodeAffinitySchedulingStrategy绑定到具体节点。Daft 把分布式物理计划切分为 Swordfish 任务分发到这些 worker 上并行执行——这也是 Daft 能在多 GPU 机器上把计算并行到 CPU 与 GPU 的原因。方式一本地单机 Ray 集群Simple Local Setup最轻量的方式是在本机启动单节点 Ray 集群pip install daft[ray] ray start --headray start --head成功后会输出类似Usage stats collection is enabled. To disable this, add --disable-usage-stats to the command that starts the cluster, or run the following command: ray disable-usage-stats before starting the cluster. Local node IP: 127.0.0.1 -------------------- Ray runtime started. -------------------- ...拿到本机 IP 与端口后通过daft.set_runner_ray把地址传给 Daft import daft daft.set_runner_ray(ray://127.0.0.1:10001) df daft.from_pydict({ ... text: [hello, world] ... }) print(df) ╭───────╮ │ text │ │ --- │ │ String │ ╞═══════╡ │ hello │ ├╌╌╌╌╌╌╌┤ │ world │ ╰───────╯ (Showing first 2 of 2 rows)默认行为如果不指定任何地址Daft 会在本机自动拉起一个本地 Ray 实例addressNone时连接或启动一个本地 Ray 实例见 set_runner_ray 文档串。对于配备了多张 GPU 的高性能单机这已经非常实用——Daft 会在 CPU 和 GPU 之间并行化执行。也可以通过环境变量DAFT_RUNNERray隐式选择 Ray runnerdaft/runners/init.py 的说明适合不想在脚本中硬编码 runner 的场景。方式二连接已有远程 Ray 集群如果已经有一个远程 Ray 集群在运行只需向set_runner_ray传入地址即可daft.set_runner_ray(addressray://url-to-mycluster)address关键字参数的完整语义与 Ray 官方ray.init一致例如ray://协议地址、auto等。注意调用set_runner_ray后 runner 配置会被锁定进程生命周期内不可再切换——tests/test_context.py 中的test_explicit_set_runner_ray等用例验证了显式设置、隐式推断get_or_infer_runner_type依次按已设置 → 检测到 Ray 集群 →DAFT_RUNNER环境变量三级策略推断 runner 类型以及各种切换限制行为。方式三使用 Ray ClientRay Client 是一种快速在远端 Ray 上运行任务并取回结果的方式import daft import ray # 注意 runtime_env 的作用参见下文 Ray Job 部分 ray.init(ray://head_node_host:10001, runtime_env{pip: [daft]}) # 启动 Ray client 并告知 Daft 使用 Ray 执行查询 # 如果 ray.init() 已经被调用过则复用现有 client daft.set_runner_ray(ray://head_node_host:10001) df daft.from_pydict({ a: [3, 2, 5, 6, 1, 4], b: [True, False, False, True, True, False] }) df df.where(df[b]).sort(df[a]) # Daft 在远端执行查询并把预览结果返回给客户端 df.collect()输出╭───────┬─────────╮ │ a ┆ b │ │ --- ┆ --- │ │ Int64 ┆ Boolean │ ╞═══════╪═════════╡ │ 1 ┆ true │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 3 ┆ true │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 6 ┆ true │ ╰───────┴─────────╯ (Showing first 3 of 3 rows)!!! warning**版本匹配要求**使用 Ray Client 运行任务时客户端与服务端的 **Daft 版本**以及 **Python 次版本号**如 3.9、3.10必须完全一致否则会出现序列化/反序列化不兼容问题。set_runner_ray还提供force_client_mode参数设为True时强制 Ray 以 client 模式运行参数说明。方式四使用 Ray Jobs推荐的生产方式Ray Jobs 相比 Ray Client 提供更强的控制力与可观测性。更重要的是你的整个代码都运行在 Ray 集群上因此不受本地机器的算力、网络、库版本和可用性限制。编写作业脚本例如wd/job.py# wd/job.py import daft def main(): # 不带任何参数调用即从 head node 连接 Ray daft.set_runner_ray() # ... 在此运行 Daft 命令 ... if __name__ __main__: main()用 Ray CLI 提交该作业CLI 通过pip install ray[default]安装ray job submit \ --working-dir wd \ --address http://head_node_host:8265 \ --runtime-env-json {pip: [daft]} \ -- python job.py!!! note--runtime-env-json {pip: [daft]} 这一 runtime env 参数的作用是**在 Ray worker 上安装 Daft**。向 worker 注入依赖还有其他替代方式如 working_dir、pip 可编辑安装、镜像预装等可按需选择。由于作业代码在集群内运行daft.set_runner_ray()可以不带参数直接连接所在集群的 head node这也是 Ray Jobs 相比客户端模式最大的运维简化。自动扩缩容scale-up 与 scale-in 机制这是 Daft 与 Ray 集成中技术含量最高的部分。当 Daft 运行在由 Ray autoscaler 管理的集群上包括 KubeRay时它会根据待执行任务pending tasks的资源需求发送 scale-up 请求。但 Ray 的 autoscaler 请求 API 是sticky的请求会粘住即使负载变空闲autoscaler 也可能一直保留之前请求的容量。Daft 因此提供了一套**可选opt-in**的机制在空闲时回收retireDaft 自管的 Flotilla worker并清除挂起的 autoscaler 请求帮助 Ray 把集群缩回去。通过 set_runner_ray 开启 scale-inimport daft daft.set_runner_ray( addressray://head_node_host:10001, downscale_enabledTrue, downscale_idle_seconds60, min_survivor_workers1, pending_release_exclude_seconds120, )对应的环境变量适合 Ray Jobs / KubeRay manifest环境变量默认值作用DAFT_AUTOSCALING_DOWNSCALE_ENABLEDfalse是否启用空闲 worker 回收scale-in。源码中1或true不区分大小写均视为开启DAFT_AUTOSCALING_DOWNSCALE_IDLE_SECONDS60worker 需空闲多久秒才成为回收候选DAFT_AUTOSCALING_MIN_SURVIVOR_WORKERS1即使空闲也必须保活的最少 worker 数防止短暂空闲把集群缩到零DAFT_AUTOSCALING_PENDING_RELEASE_EXCLUDE_SECONDS120被回收 worker ID 的黑名单宽限 TTL秒防止 autoscaler 立即重新拉起同规格节点造成抖动set_runner_ray()的参数与这些环境变量是等价的两条入口Python 包装层直接把传参写入同名环境变量daft/runners/init.py再由 Rust 侧统一读取注释里明确说明这样设计是为了让配置经由环境变量传递到 Rust 调度器/worker 管理器而无需在整个技术栈中层层穿参。源码深潜retire_idle_workers 的完整决策链回收逻辑全部实现在 worker_manager.rs 的retire_idle_workers中其决策顺序为读开关先读DAFT_AUTOSCALING_DOWNSCALE_ENABLED未开启直接返回 0读取保活下限DAFT_AUTOSCALING_MIN_SURVIVOR_WORKERS默认 1集群彻底空闲时的强制清扫若处于最终关停周期force_all_when_cluster_idle会无条件调用 flotilla 的clear_autoscaling_requests()见 flotilla.py内部即request_resources(bundles[])清空所有挂起的 autoscaler 需求且绕过空闲时长阈值与保活下限允许回收全部 workerscale-up 保护若正处于活跃的扩容窗口skip_due_to_pending_scale_up本轮跳过回收避免刚发给 Ray 的扩容需求被自己抵消候选筛选遍历所有 workerhead node 直接豁免通过 get_head_node_id 识别它读取 Ray 内部资源键node:__internal_head__只有处于空闲状态且空闲时长 ≥downscale_idle_seconds的 worker 才成为候选按空闲时长降序选取候选按空闲最久优先排序取总 worker 数 − min_survivor_workers个执行释放释放的 worker ID 写入pending_release_blacklist并打上时间戳在 TTL 内被 refresh_workers 排除在 worker 发现之外清除 autoscaler 需求再次调用clear_autoscaling_requests()让 Ray autoscaler 有机会真正缩容。这套设计解释了文档中sticky request问题的完整解法既回收 worker又清需求还通过黑名单 TTL 防止回收—重拉的振荡。scale-up 侧gradual 与 bisect 两种策略set_runner_ray还有两个文档未展开、但源码中完整实现的扩参数autoscale_strategy取gradual默认或bisectautoscale_bisect_timeout_secs仅 bisect 策略使用等待集群扩容的超时秒数必须大于 0默认 30。gradual渐进式策略try_autoscale_gradual基于一个朴素但稳健的观察Daft 无法预知集群的扩容上限如 KubeRay 的maxReplicas而ray.autoscaler.sdk.request_resources是异步的、每次调用会整体替换非累加当前需求且 Ray autoscaler 每约 5 秒对账一次可经AUTOSCALER_UPDATE_INTERVAL_S调整Daft 读取同一变量来对齐自己的请求节奏。若一次性请求超过集群上限autoscaler 会整体拒绝而非部分满足。因此 gradual 策略维护一个高水位每个 autoscaler 对账周期只比上一次请求多请求一个 bundle逐步逼近集群真实容量上限高水位同时被当前集群实际资源抬升floor冷启动时第一周期即可跳过已有容量直接请求增量。bisect对半收敛策略try_autoscale_bisect则采用二分思想首次一次性请求全部待执行需求若autoscale_bisect_timeout_secs内集群没有增长worker 集合严格扩张或 CPU/GPU/内存任一维度上升见cluster_capacity_grew判定就认定请求被拒并把请求量减半下限为单个 bundle 的需求循环往复以 O(log N) 步收敛到集群真实容量若集群确实增长则贪心地重新请求全部剩余需求。值得注意的是增长判定中对worker 替换的区分worker_set_grew 要求新 worker 集合是旧集合的严格超集因此一个旧 worker 被同容量新 worker 替换不算扩容——对应的单元测试equal_capacity_worker_replacement_is_not_growth等用例在 worker_manager.rs 的 tests 模块 中逐条验证了这些边界。参数校验也有测试兜底tests/test_context.py 验证了autoscale_bisect_timeout_secs0会被拒绝。此外还有worker_startup_timeout参数对应环境变量DAFT_RAY_WORKER_STARTUP_TIMEOUT控制 Ray worker actor 上报地址的启动超时供慢启动环境调优。扩容与回收如何互相配合从源码看两处精巧的联动worker_manager.rs发出新扩容请求时会立即清空pending_release_blacklist并强制刷新 worker 列表——这样刚被回收节点上的 worker 可以立刻重建新供应的节点也能被快速观察到回收发生期间处于活跃 scale-up 的周期会整体跳过回收前述第 4 步避免一边要扩容、一边缩 worker的自相矛盾。参数速查与适用前提汇总set_runner_ray的完整签名daft/runners/init.py参数类型默认说明addressstr \| NoneNoneRay 集群地址None时连接或启动本地 Ray 实例noop_if_initializedboolFalseRay 已初始化时跳过初始化测试与 notebook 环境常用如 tests/extensions/test_extension_runtime.py 的用法force_client_modeboolFalse强制 Ray 以 client 模式运行downscale_enabledbool \| None回退环境变量默认关闭开启空闲 worker 回收scale-indownscale_idle_secondsint \| None60worker 空闲多久可被回收必须 ≥ 0min_survivor_workersint \| None1保活的最少 worker 数必须 ≥ 0pending_release_exclude_secondsint \| None120被回收 worker ID 的黑名单 TTL必须 ≥ 0worker_startup_timeoutint \| None见 flotilla 默认值worker 上报地址的启动超时秒可用DAFT_RAY_WORKER_STARTUP_TIMEOUT覆盖autoscale_strategystr \| Nonegradual扩容策略gradual或bisectautoscale_bisect_timeout_secsint \| None30bisect 策略的扩容观察窗口必须 0适用前提与限制扩缩容特性依赖 Ray autoscaler 管理的集群含 KubeRay在固定规模的手动集群上 scale-in 只会回收 Daft 自管的 worker不会改变集群节点数Ray Client 模式要求客户端与服务端 Daft 版本和 Python 次版本号严格一致set_runner_ray之后 runner 被锁定同一进程内不可再切换回 native runner若使用 Ray Jobs 部署到 Kubernetes可参考仓库中的 Helm chart k8s/charts/quickstart 以及 Kubernetes 部署文档其中环境变量入口DAFT_AUTOSCALING_*正是为这类 manifest 场景设计的——在环境变量里配置比改脚本更自然。小结Daft 的 Ray 集成提供了从轻到重的四条路径本机ray start --head快速上手、远程集群地址直连、Ray Client 快速开发、Ray Jobs 生产提交在此之上downscale_*参数族与autoscale_strategy让 Daft 能够与 Ray autoscaler 双向协商容量——gradual 策略保守稳健bisect 策略快速收敛scale-in 则通过空闲回收 需求清零 回收黑名单三重机制解决 Ray autoscaler 请求粘住不缩的固有痛点。理解 worker_manager.rs 中这套纯 Rust 侧的决策逻辑是调优 Daft 大规模 Ray 部署的关键。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考