尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

在 Flyte 工作流中运行 Daft:Ray 集群接入与多模态数据处理实战

在 Flyte 工作流中运行 Daft:Ray 集群接入与多模态数据处理实战 在 Flyte 工作流中运行 DaftRay 集群接入与多模态数据处理实战【免费下载链接】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本篇技术指南讲解如何在 Flyte 工作流引擎中运行 Daft重点介绍两种典型用法在普通 Flyte 任务中像本地机器一样直接使用 Daft以及在 Ray 类型 Flyte 任务中通过daft.set_runner_ray()接入 Ray 集群进行分布式执行。读完本文你将掌握 Flyte 任务的 Daft 环境搭建、Ray Runner 的接入原理以及一个完整的S3 读取 → 文本过滤 → 图片下载缩放 → Parquet 写出多模态数据流水线示例。Flyte 与 Daft 的结合方式Flyte 是流行的开源工作流引擎其 Ray 插件允许任务运行在 Ray 集群上。Daft 与该工作流引擎及其 Ray 插件配合良好集成方式非常轻量Daft 不需要特殊的 Flyte API只需要在任务函数中正确选择执行 Runner。Daft 默认使用 Native Runner多线程本地执行而通过daft.set_runner_ray()可以切换到 Ray Runner从而把 DataFrame 的计算分发到 Ray 集群上执行。在 Flyte 场景下有两种常见形态场景使用方式说明普通 Flyte 任务直接使用 Daft API与在本地机器上使用 Daft 完全一致Ray Flyte 任务先调用daft.set_runner_ray()复用 Flyte 已初始化的 Ray 集群连接在普通 Flyte 任务中使用 Daft在 Flyte 任务中运行 Daft 非常简单——像在本地机器上使用 Daft 一样使用即可。Flyte 任务是一个被调度执行的 Python 函数Daft 的惰性求值、查询优化与执行机制在其中照常工作from flytekit import task task() def my_daft_task(limit: int) - list[str]: import daft df daft.read_parquet(s3://your-bucket/data.parquet) result df.where(df[col] 0).limit(limit).collect() return result.to_pydict()[id]此时 Daft 默认使用 Native Runner任务内部多线程并行执行无需额外配置。在 Ray Flyte 任务中使用 Daft当任务运行在 Ray 集群上时通过 Flyte 的 Ray 插件Daft 会自动感知环境并接入 Ray。核心要点是如果 Ray 已经通过ray.init()初始化Daft 会复用已有的 Ray 集群连接无需再次初始化。从源码实现来看RayRunner 的初始化逻辑 正是如此设计的若ray.is_initialized()返回 True即 Flyte 的 Ray 插件已经初始化了 RayDaft 直接复用现有 Ray context即使传入了address也仅打印警告并忽略若 Ray 尚未初始化则调用ray.init(addressaddress)建立连接。因此在 Ray Flyte 任务中只需一行代码daft.set_runner_ray()Daft 会遵循以下判定顺序对应 get_or_infer_runner_type 的推断策略若 Runner 已被显式设置直接使用检测是否运行在 Ray 集群中Flyte Ray 任务即属于此情况回退到DAFT_RUNNER环境变量取值为native或ray。set_runner_ray 的完整参数除了无参调用set_runner_ray 还支持丰富的参数用于精细控制 Ray 连接与集群扩缩容行为参数默认值说明addressNoneRay 集群地址为None时连接或启动本地 Ray 实例。注意以ray://前缀连接会走 Ray Client 模式可能影响性能源码中对此有显式警告见 ray_runner.pynoop_if_initializedFalse若 Ray 已运行则跳过初始化force_client_modeFalse强制 Ray 以客户端模式运行downscale_enabled环境变量DAFT_AUTOSCALING_DOWNSCALE_ENABLED默认False是否启用空闲 worker 缩容scale-indownscale_idle_seconds环境变量DAFT_AUTOSCALING_DOWNSCALE_IDLE_SECONDS默认60worker 空闲多久后可被回收min_survivor_workers环境变量DAFT_AUTOSCALING_MIN_SURVIVOR_WORKERS默认1即使空闲也保留的最少 worker 数pending_release_exclude_seconds环境变量DAFT_AUTOSCALING_PENDING_RELEASE_EXCLUDE_SECONDS默认120刚释放 worker 的 TTL 宽限期避免被立即重新拉起worker_startup_timeout环境变量DAFT_RAY_WORKER_STARTUP_TIMEOUT启动时等待 Ray worker actor 上报地址的超时秒数autoscale_strategygradual扩容策略gradual逐 bundle 逐步加码或bisect一次性请求全部需求、被拒后对半缩减以加速收敛autoscale_bisect_timeout_secs30仅在bisect策略下使用等待集群扩容的超时秒数扩缩容相关参数最终会写回对应的DAFT_AUTOSCALING_*环境变量以便传递到 Rust 层的调度器与 worker 管理器见 daft/runners/init.py 的注释说明。两种 Ray 接入形态对比根据 Ray 运行文档 与源码Ray Runner 的接入主要有两种形态本地单机 Ray本地执行ray start --head启动集群再通过daft.set_runner_ray(ray://127.0.0.1:10001)连接。对于拥有多 CPU / 多 GPU 的机器如 AWS P3 实例Daft 可以在本机 CPU 与 GPU 上并行执行。远程 Ray 集群传入集群 head 节点地址即可如daft.set_runner_ray(addressray://url-to-mycluster)。在 Flyte Ray 任务场景下通常无需手动指定地址——Flyte 的 Ray 插件负责集群编排Daft 复用已初始化的 Ray context 即可。完整示例Flyte 任务中的图片处理流水线仓库中的 tutorials/flyte/app.py 提供了一个可运行的完整示例在 Flyte 任务中读取 LAION 数据集的 Parquet 文件过滤文本含 darkness 的行下载对应图片并缩放到 32×32最后写出 Parquet。任务与工作流定义from flytekit import current_context, task, workflow from flytekitplugins.ray import ( # noqa RayJobConfig, WorkerNodeConfig, ) import daft task() def produce_resized_image_dataset(limit: int) - list[str]: # NOTE: Use Ray Runner: # # 1. If Ray connection has been set up already by Flyte, # it will use the initialized Ray cluster connection # 2. Otherwise, it will create a local Ray cluster on the Flyte task # daft.set_runner_ray() written_df daft_notebook_code(limit) # Return the number of rows written return written_df.to_pydict()[file_path] workflow() def wf(limit: int): produce_resized_image_dataset(limitlimit)代码注释清晰说明了daft.set_runner_ray()的两种行为Flyte 已建立 Ray 连接则复用否则在任务内创建本地 Ray 集群——这正是前面源码分析中ray.is_initialized()分支的体现。文件顶部注释掉的RayJobConfig/WorkerNodeConfig展示了如何将任务声明为 Ray 任务并指定 worker 组# task( # task_configRayJobConfig( # worker_node_config[ # WorkerNodeConfig( # group_nameray-group, replicas1, # ) # ], # ) # )按需取消注释即可让 Flyte 为该任务编排一个名为ray-group、含 1 个副本的 Ray worker 组。数据流水线核心逻辑def daft_notebook_code(limit: int): # 以匿名模式访问 AWS S3 IO_CONFIG daft.io.IOConfig( s3daft.io.S3Config(anonymousTrue, region_nameus-west-2) ) PARQUET_PATH s3://daft-oss-public-data/tutorials/laion-parquet/train-00000-of-00001-6f24a7497df494ae.parquet parquet_df daft.read_parquet(PARQUET_PATH, io_configIO_CONFIG) parquet_df parquet_df.select( parquet_df[URL], parquet_df[TEXT], parquet_df[AESTHETIC_SCORE] ) # 过滤出文本包含 darkness 的样本 filtered_df parquet_df.where(parquet_df[TEXT].contains(darkness)) # 下载图片并按 URL 解码 filtered_df filtered_df.with_column( image, filtered_df[URL].download(on_errornull).decode_image(), ) # 缩放到 32x32 filtered_df filtered_df.with_column( resized_image, filtered_df[image].resize(32, 32) ) # 写出 Parquet生产环境应写云存储示例写入本地临时文件 written_df ( filtered_df.select(URL, TEXT, resized_image) .limit(limit) .write_parquet( fmy-s3-bucket/{current_context().execution_id.name}/resized_images.parquet ) ) return written_df流水线各步骤要点匿名 S3 读取daft.io.S3Config(anonymousTrue, region_nameus-west-2)无需凭证即可访问公开数据集on_errornull单个 URL 下载失败时将该值置为 null 而不是中断整个任务提升批量处理鲁棒性注意 notebook 中decode_image也带有on_errornull参数与app.py的用法略有差异两者都是合法的错误处理方式惰性求值以上所有操作都是构建逻辑计划真正执行发生在write_parquet返回后结果回传任务返回写出的文件路径列表供 Flyte 下游任务消费。任务运行时的执行上下文示例通过current_context().execution_id.name获取 Flyte 执行 ID 用于构造输出路径保证每次运行写入独立目录避免任务重试或并行运行时的路径冲突。本地运行环境搭建要本地运行该示例需要先完成 Flyte 本地环境包括 Ray 插件的搭建具体步骤参考 Flyte 官方 Ray 示例仓库。仓库中配套的 tutorials/flyte/Dockerfile 给出了镜像构建的最小依赖FROM ghcr.io/flyteorg/flytekit:py3.9-1.7.0 USER root RUN pip install daft flytekitplugins-ray即以 Flyte 官方 kit 镜像为基础额外安装daft与flytekitplugins-ray两个包即可。若直接使用本地 Python 环境安装命令等价于pip install daft[ray] flytekit flytekitplugins-raydaft[ray]的 extra 依赖会一并安装 Ray具体依赖声明可参考 pyproject.toml 中的依赖配置。使用 notebook 先行验证流水线仓库中的 tutorials/flyte/notebook.ipynb 是同一流水线的交互式版本适合在本地先用 Daft 验证数据处理逻辑再将其封装进 Flyte 任务。notebook 按步骤演示了配置匿名 S3 IOConfig 并读取 Parquetselect选择URL/TEXT/AESTHETIC_SCORE三列collect()/show(5)预览数据where(TEXT.contains(darkness))文本过滤download(on_errornull).decode_image(on_errornull)下载并解码图片resize(32, 32)缩放图片write_parquet写出结果。注意 notebook 使用的数据路径为s3://daft-public-data/...与app.py中的s3://daft-oss-public-data/...不同以实际可访问的公开数据集为准。该 notebook 本质上与 quickstart 教程 的图像处理思路一致可对照学习 Daft 的惰性执行与多模态数据处理 API。常见问题与注意事项无需手动ray.init()在 Ray Flyte 任务中不要重复初始化 Ray。Daft 检测到已初始化的 Ray context 时会复用重复初始化可能导致连接混乱。ray://前缀的性能影响若显式传入ray://地址Daft 会走 Ray Client 模式并打印性能警告见 ray_runner.py。在 Flyte Ray 任务内部运行时通常不需要指定地址Daft 自动复用集群连接即可。Runner 一经设置不可中途变更根据 get_or_create_runner 的文档说明Runner 在进程生命周期内一旦锁定便无法更换因此应在任务函数开头尽早调用set_runner_ray()。临时文件系统注意示例中将结果写入本地临时路径。在容器化的 Flyte 任务中本地文件随任务结束而消失生产环境应写回 S3、GCS 等持久化存储。任务返回类型示例任务返回list[str]写出文件的路径列表Flyte 会将其作为任务输出进行传递与记录RayJobConfig声明 Ray 任务时同理可返回任意可序列化类型。小结Flyte 与 Daft 的集成的核心就两步普通任务直接使用 Daft APIRay 任务在函数开头调用一次daft.set_runner_ray()。Daft 会智能复用 Flyte Ray 插件已初始化的集群连接否则在任务内自建本地 Ray 集群。结合 tutorials/flyte/app.py 的完整示例与 notebook.ipynb 的交互式验证流程你可以快速把本地验证过的 Daft 多模态流水线读 S3、过滤、下载图片、缩放、写 Parquet平滑迁移到 Flyte 的分布式调度体系中。【免费下载链接】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),仅供参考
返回列表