ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Metaflow R 流程云端化实战:用 `--with batch` 在 AWS Batch 上运行 MovieStatsFlow 而不修改一行代码

Metaflow R 流程云端化实战:用 `--with batch` 在 AWS Batch 上运行 MovieStatsFlow 而不修改一行代码 MLOps工作流自动化数据工程【免费下载链接】metaflowBuild, Manage and Deploy AI/ML Systems项目地址https://gitcode.com/gh_mirrors/me/metaflow点击查看免费下载本篇指南围绕仓库中 Episode 05-statistics-redux 教程 展开讲解如何把 Episode 02 中完全本地执行的 R 统计流程stats.R原封不动地搬到 AWS Batch 云端执行只追加一个--with batch命令行参数Metaflow 就会把所有步骤调度到远程计算资源上数据 artifact 自动落盘 AWS S3随后你可以在本地 RStudio 中通过 Metaflow Client 重新读取这些云端数据并完成分析。读完本文你将掌握 Metaflow R API 的云端运行命令、--max-workers与--package-suffixes的用法以及云端 artifact 的本地回读流程。教程背景从 Episode 02 到 Episode 05Episode 05-statistics-redux 是 Episode 02stats.R电影类型统计流程的云端重制版。两者共享完全相同的流程代码start读取 movies.csv 并枚举全部电影类型compute_stats以foreach方式按类型并行计算票房的中位数median与均值mean最后join汇总结果。完整流程见 R/inst/tutorials/02-statistics/stats.Rmetaflow(MovieStatsFlow) %% step(step start, r_function start, next_step compute_stats, foreach genres) %% step(step compute_stats, r_function compute_stats, next_step join) %% step(step join, r_function join, next_step end, join TRUE) %% step(step end) %% run()Episode 05 的关键主张是用 Metaflow 做云端扩展不需要改动任何业务代码。Episode 02 中每个步骤默认在本地进程中顺序/并行执行Episode 05 只是为同一条命令追加--with batchMetaflow 便自动将流程图的每个步骤打包、提交到 AWS Batch 上运行而compute_stats的foreach分支会在云端按--max-workers指定的并行度展开。开始之前准备 AWS sandbox 环境教程要求先配置好可用的 Metaflow AWS sandbox。其要点是环境已配置好 AWS 凭据与默认 region默认元数据后端metadata provider指向 Metaflow service默认数据存储datastore为 S3对应 run.R 中datastore: local (default) or s3的说明AWS Batch 计算环境、作业队列以及 S3 访问 IAM 角色均已就绪。只有在这种云端就绪的环境里--with batch才能把任务提交到真实 AWS Batch 队列若在纯本地环境使用该参数Metaflow 会因缺少队列配置而报错。batch 相关默认值在 R/R/decorators-aws.R 中有明确来源队列取自环境变量METAFLOW_BATCH_JOB_QUEUES3 访问角色取自METAFLOW_ECS_S3_ACCESS_IAM_ROLEFargate 执行角色取自METAFLOW_ECS_FARGATE_EXECUTION_ROLE。在终端执行云端流程进入 Episode 02 的教程目录运行cd tutorials/02-statistics/ Rscript stats.R --package-suffixes.R,.csv run --with batch --max-workers 4命令拆解如下片段作用Rscript stats.R以非交互方式执行流程脚本stats.R末尾的run()会接管剩余参数--package-suffixes.R,.csv把当前目录下的.R与.csv文件一并打包进代码包code package随任务上传到云端run执行一次完整运行区别于resume等子命令--with batch给流程整体挂上 batch 装饰器所有步骤调度到 AWS Batch 执行--max-workers 4限制同时运行的并行任务数为 4即compute_stats的 foreach 分支最多 4 个并发--package-suffixes是本例中容易忽略但很关键的一环stats.R的start步骤执行read.csv(./movies.csv, stringsAsFactorsFALSE)依赖本地 CSV 文件。--package-suffixes.R,.csv保证movies.csv进入代码包随任务一起分发给云端各 worker因此云端步骤也能读到同样的数据文件。这正是文档强调甚至本地 CSV 文件也被存起来的原因——不仅数据 artifact 落盘 S3整个代码包含.R脚本与.csv输入同样被持久化。在 RStudio 中执行云端流程如果使用 RStudio 交互环境把stats.R最后一行run()替换为run(batchTRUE, max_workers4, package_suffixes.R,.csv,)然后执行source(stats.R)。从源码看这一调用与命令行是完全等价的。run()定义于 R/R/run.R它会先把流程对象序列化为临时 RDS 文件再调用run_cmd()把 R 命名参数翻译成 CLI 参数batchTRUE被转换为--with batch见run_cmd中if (is.logical(flags$batch)) { if (flags$batch) batch - --with batch }max_workers4被转换为--max-workers4package_suffixes被拼接为--package-suffixes.R,.csv,逗号分隔见 R/R/run.R 第 154-158 行。最终命令通过system()以Rscript run.R --flowRDS... --no-pylint --package-suffixes... --with batch run --max-workers4的形式启动R/inst/run.R 负责反序列化 RDS、把各步骤 R 函数注入全局环境再经由reticulate调用 Python 端metaflow$R$run真正执行流程。可见run()只是一个参数桥交互式与命令行两种方式殊途同归。--max-workers的作用与语义--max-workers限制整个流程中同时运行的 task 数量。在本例中genres有多个取值compute_stats会按类型拆出多个并行分支--max-workers 4即保证这些分支最多 4 个并发执行避免一次性拉起过多 AWS Batch 任务、控制云成本与配额。在 Python 侧该参数由 metaflow/cli_components/run_cmds.py 第 147 行附近定义并传入运行逻辑当流程被部署到 Argo Workflows 时它映射为工作流的parallelism上限见 metaflow/plugins/argo/argo_workflows.py部署到 AWS Step Functions 时映射为max_concurrency见 metaflow/plugins/aws/step_functions/step_functions.py。需要说明--max-workers只控制并发 task 数量的上限不改变流程图的逻辑结构。本地检查云端结果通过 Metaflow Client 读取 S3 artifact云端运行结束后打开 R/inst/tutorials/02-statistics/stats.Rmd在本地 RStudio 中重新运行其中的 cell即可通过 Metaflow Client 拿到最新一次成功运行的statsartifact 并绘图suppressPackageStartupMessages(library(metaflow)) message(Current metadata provider: , get_metadata()) message(Current namespace: , get_namespace()) flow - flow_client$new(MovieStatsFlow) run_id - flow$latest_successful_run run - run_client$new(flow, run_id) df - run$artifact(stats) print(head(df)) df - df[order(df$median, decreasing TRUE), ] print(head(df)) barplot(df$median[1:5], names.argdf$genres[1:5])这段代码揭示了云端数据回读的完整链路flow_client$new(MovieStatsFlow)创建一个指向已有流程的客户端R/R/flow_client.Rlatest_successful_run自动定位最近一次成功运行的 run_idrun_client$new(flow, run_id)封装该次运行R/R/run_client.Rrun$artifact(stats)从运行结果中取出join步骤产出的数据框并反序列化为本地 R 对象数据实际存储在 S3 数据存储中Metaflow Client 按需从 S3 拉取并反序列化因此流程在云端跑、分析在本地做成为可能。stats.Rmd开头还打印当前 metadata provider 与 namespace用于确认客户端确实连到了 sandbox 的元数据服务而不是本地默认元数据。底层原理Metaflow 如何把 R 步骤送上 AWS Batch从仓库实现看云端调度的核心在 metaflow/plugins/aws/batch/batch.pyBatch类共 588 行每个步骤执行时Metaflow 把代码包含 RDS 序列化的流程对象、--package-suffixes指定的附属文件上传到 S3 数据存储通过 AWS Batch 提交一个容器任务容器内先恢复代码包再以Rscript重新进入 R/inst/run.R 的运行时入口执行对应步骤的 R 函数步骤的日志经由 mflog 机制bash_capture_logs、tail_logs收集回元数据服务供本地查询步骤产出的 artifactself$df、self$median、self$stats等写入数据存储因此run$artifact(stats)才能跨机器取回。这也解释了教程开篇那句话的完整含义数据被存储在 AWS S3 中因此你可以从任何地方访问它们——流程代码、输入 CSV、中间与最终 artifact 全部落在云存储上本地客户端只是数据的又一个消费者。延伸进一步控制云上资源--with batch给流程整体挂上 batch 装饰器而如果你需要对单个步骤定制 CPU、内存等资源可以在 R 流程中对特定 step 使用batch(...)装饰器。其完整签名见 R/R/decorators-aws.Rcpu所需 CPU 数默认1gpu所需 GPU 数默认0memory所需内存MB默认4096image容器镜像缺省时使用合适的 Rocker Docker 镜像queue提交作业的 AWS Batch 队列默认取METAFLOW_BATCH_JOB_QUEUEiam_roleBatch 访问 S3 的 IAM 角色默认取METAFLOW_ECS_S3_ACCESS_IAM_ROLEexecution_role触发 Fargate 任务的执行角色默认取METAFLOW_ECS_FARGATE_EXECUTION_ROLEshared_memory、max_swap、swappiness容器内存与交换行为调优。当同时存在resources与batch装饰器时Metaflow 取所有装饰器中的最大值作为最终资源请求。若流程中的大内存步骤需要 60GB 内存可仿照该文件文档示例把batch(memory60000, cpu1)挂在对应 step 上——这属于 Episode 05 之外的进阶话题但和--with batch属于同一套云端资源控制体系。小结Episode 05 用最少的动作展示了 Metaflow 的可移植性核心能力零代码迁移同一份 stats.R 流程加一个--with batch即从本地扩展为 AWS Batch 云端执行并行度可控--max-workersR 侧max_workers限制云端并发任务数兼顾性能与成本依赖文件随行--package-suffixes.R,.csv确保云端 worker 也能读取movies.csv等输入文件数据云端持久化artifact 与代码包存入 S3本地 RStudio 借助flow_client/run_client/run$artifact()即可回读并完成分析。这套本地开发、云端运行、本地分析的工作流是 Metaflow R API 在真实云环境中的标准用法也是从原型走向生产部署如 Argo Workflows、Step Functions之前最平滑的一步。赞分享MLOps工作流自动化数据工程【免费下载链接】metaflowBuild, Manage and Deploy AI/ML Systems项目地址https://gitcode.com/gh_mirrors/me/metaflow点击查看免费下载相关推荐Apache Airflow AWS Batch Executor 实战指南以 Amazon Batch 弹性运行工作流Apache Airflow AWS Batch Executor 实战指南以 Amazon Batch 弹性运行工作流 Apache Airflow 的 A后端任务调度工作流自动化数据编排批处理数据工程流程编排AWS CLI 实战指南使用 aws batch terminate-job 终止 AWS Batch 作业AWS CLI 实战指南使用 aws batch terminate job 终止 AWS Batch 作业 导读 本文以 AWS CLI 官方示例文档 te开发工具云原生运维DyberPet让桌面拥有生命感的数字伙伴框架DyberPet让桌面拥有生命感的数字伙伴框架 你是否曾觉得电脑桌面太过冰冷缺少一丝温暖与陪伴每天面对单调的工作界面是否渴望有一个会呼吸、会互动的小生命开发工具云原生运维创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表