- 后端
- 任务调度
- 工作流自动化
- 云原生
- MLOps
- 微服务
【免费下载链接】flyte
Dynamic, resilient AI orchestration. Coordinate data, models, and compute as you build AI workflows.
Flyte CoPilot 是 Flyte 面向"任意容器"场景的关键基础设施:它通过一个理解 FlyteIDL 元数据格式的flyte-copilot二进制,把 Flyte 的输入输出协议翻译成普通的本地文件系统操作,从而让完全不了解 Flyte SDK 的容器也能作为 Flyte 任务运行。本文以 flytecopilot/README.md 为骨架,结合仓库内 CLI 实现、数据层源码与 flyteidl2/core/tasks.proto 协议定义,讲透它的两种运行模式、容器生命周期编排思路、全部命令行参数以及底层数据搬运原理。
核心思想:把 Flyte 协议翻译成"文件系统协议"
传统 Flyte 任务容器内通常需要安装 flytekit(Python)等 SDK 才能与 Flyte 平台交换输入输出。CoPilot 的目标是消除这一依赖:让任意镜像、任意语言写的容器,只需要读写本地文件系统就能参与 Flyte 任务。
从 flyteidl2/core/tasks.proto 中DataLoadingConfig的注释可以确认这一设计意图:
This configuration allows executing raw containers in Flyte using the Flyte CoPilot system. Flyte CoPilot, eliminates the needs of flytekit or sdk inside the container. Any inputs required by the users container are side-loaded in the input_path. Any outputs generated by the user container - within output_path are automatically uploaded.
也就是说,CoPilot 负责两件事:
- 输入侧加载(side-load):把远程存储中的输入元数据(Flyte Metadata Format,即
core.LiteralMap)和实际数据下载到容器可见的本地路径; - 输出侧上传(auto-upload):主容器结束后,把约定路径下生成的输出文件(尤其是元数据
outputs.pb)上传回远程存储。
这一切由同一个二进制flyte-copilot的两种模式完成。从 flytecopilot/main.go 看,二进制入口会构造根命令NewDataCommand()并挂载download与sidecar两个子命令(见 flytecopilot/cmd/root.go),根命令的定位描述是"flytedata is a simple go binary that can be used to retrieve and upload data from/to remote stow store to local disk",与 README 中"2 种模式"的说法一致。
两种运行模式总览
| 模式 | CLI 子命令 | 部署形态 | 职责 |
|---|---|---|---|
| Downloader | flyte-copilot download | K8s init container | 下载元数据与数据到共享卷 |
| Sidecar | flyte-copilot sidecar | K8s sidecar 容器 | 监控主容器生命周期并上传输出 |
需要说明的是:README 以 "Downloader" 指代该模式,而源码中 cobra 子命令的实际名称是download(见 flytecopilot/cmd/download.go),二者是同一功能的两种叫法。
Mode: Downloader —— 在 init container 中预置输入
基本用法
$ flyte-copilot download在 Kubernetes 中,download通常作为init container运行,并将下载卷(shared volume)挂载到主容器:
- init container 先于主容器启动并完成数据下载;
- 下载完成后,主容器启动时,输入数据已全部就位在本地路径上;
- 这保证了主容器永远看不到"半成品"输入。
从 flytecopilot/cmd/download.go 的实现看,Download()会依次校验to-output-prefix、to-local-dir、from-remote三个必填参数,再校验format、download-mode、file-input-layout的合法性,最终调用data.NewDownloader(...).DownloadInputs(...)完成递归下载。
完整命令行参数
| 参数 | 短选项 | 默认值 | 说明 |
|---|---|---|---|
--from-remote | -f | 无(必填) | 远程存储中输入元数据(inputs.pb)的路径/键 |
--to-output-prefix | 无 | 无(必填) | 输出元数据前缀,主要用于写error.pb |
--to-local-dir | -o | 无(必填) | 本地下载目录 |
--format | -m | JSON | 元数据编码格式,可选JSON/YAML/PROTO |
--download-mode | -d | DOWNLOAD_EAGER | 下载模式(见下文 IOStrategy) |
--file-input-layout | 无 | DIRECT | File/list[File] 输入在磁盘上的布局 |
--timeout | -t | 1h | 下载允许的最大时长 |
--input-interface | -i | 无 | base64 编码的core.VariableMap(输入接口 protobuf) |
所有取值校验都直接对照flyteidl2/core/tasks.proto中的枚举实现(见 flytecopilot/cmd/download.go),取值不合法时会返回带全部可选值的明确报错,例如incorrect input download format specified, given [xxx], possible values [JSON YAML PROTO]。
下载后的目录布局
从DataLoadingConfig.input_path的注释(flyteidl2/core/tasks.proto)可以归纳出输入目录的标准形态,例如输入接口为(x: int, y: blob, z: multipart_blob)、输入路径为/var/flyte/inputs:
/var/flyte/inputs/ ├── inputs.json # 元数据摘要文件(依 format 可为 .pb/.json/.yaml) ├── inputs.pb # 始终写入的 protobuf 版 LiteralMap ├── x # 整数 x 的字符串形式 ├── y # Blob y 的二进制内容 └── z/... # multipart blob z 是一个目录,内部为各 part从 flytecopilot/data/download.go 的DownloadInputs()可以看到:下载器总是写出inputs.pb,当format=JSON时额外写出inputs.json,当format=YAML时额外写出inputs.yaml;远程引用在元数据中会被替换为本地文件系统路径(handleBlob返回toPath并改写 scalar 的 URI)。
Mode: Sidecar —— 守护主容器生命周期并回收输出
五步目标
README 明确了 sidecar 模式的核心目标流程:
- 识别主容器(identify the main container)
- 等待主容器启动(wait for the main container to start up)
- 等待主容器退出(wait for the main container to exit)
- 把数据(尤其是元数据)复制到远程存储(copy the data to remote store)
- 退出(exit)
$ flyte-copilot sidecar从 flytecopilot/cmd/sidecar.go 的Sidecar()看,流程被实现为uploader()+ 错误兜底:正常路径下等待容器退出后执行RecursiveUpload;若上传过程中出现RawContainerError(主容器以 Flyte 错误文档形式失败),则直接上传该错误文档;其他失败则统一上传OutputUploadFailed错误文档。
完整命令行参数
| 参数 | 短选项 | 默认值 | 说明 |
|---|---|---|---|
--to-output-prefix | -o | 无(必填) | 输出元数据在 stow store 中的远程前缀 |
--to-raw-output | -x | 无 | 原始输出数据的远程前缀(沙箱目录) |
--from-local-dir | -f | 无 | 主容器输出所在的本地目录 |
--format | -m | JSON | 原始/结构化类型的输出元数据格式 |
--upload-mode | -u | UPLOAD_ON_EXIT | 上传时机(见下文 IOStrategy) |
--meta-output-name | 无 | outputs.pb | 成功执行时输出元数据文件的键名 |
--timeout | -t | 1h | 上传允许的最大时长(旧名--start-timeout已废弃但保留兼容) |
--interface | -i | 无 | base64 编码的core.TypedInterface,声明输出接口 |
--start-watcher-type | 无 | signal | 等待容器"启动"的 watcher 类型 |
--exit-watcher-type | 无 | signal | 等待容器"退出"的 watcher 类型 |
输出接口驱动的上传
Sidecar 上传并非盲目拷贝整个目录,而是严格按TypedInterface的输出变量逐个处理(flytecopilot/cmd/sidecar.go):
- 未提供
--interface或输出接口为空时,直接按 Void 输出处理、立即退出; - 每个输出变量按其类型分派:简单类型(
SimpleType)走handleSimpleType,Blob 类型走handleBlobType; - 简单类型文件有1024 字节的上限校验(
maxPrimitiveSize,见 flytecopilot/data/upload.go),Blob 类型的目录会递归遍历并并发上传每个 part; - 上传完成后把所有
core.Literal聚合为LiteralMap并以 protobuf 写入outputs.pb(--meta-output-name可改名)。
值得注意的是,源码中定义了三个约定文件名常量(flytecopilot/cmd/sidecar.go):
const ( StartFile = "_START" SuccessFile = "_SUCCESS" ErrorFile = "_ERROR" )其中_ERROR是主容器输出目录中的错误文件,Uploader 在开始上传前会先检查它:若内容可反序列化为core.ErrorDocument且包含 code/message,则抛出RawContainerError让平台拿到结构化错误;否则按普通文本错误处理(见 flytecopilot/data/upload.go 与RecursiveUpload开头部分)。
如何识别主容器与感知退出:三种方案的演进
README 中保留了"Raw notes"(原始设计笔记),对比了三种识别主容器、等待其退出的方案,这对理解 watcher 抽象很有价值:
方案 1:轮询 Kube API
poll Kubeapi. Works perfectly fine, but too much load on kubeapi
通过持续调用 Kubernetes API 查询 Pod 状态。功能上可行,但会给 API Server 带来过大负载,被否决。
方案 2:约定_SUCCESS文件协议
Main container will exit and write a _SUCCESS file to a known location
主容器正常退出时向约定位置写一个_SUCCESS文件,sidecar 据此判定完成。缺点也很明显:发生 OOM 或随机退出时不会写该文件,uploader 会被卡住。README 提出的缓解思路是引入超时,或在主容器异常退出时由 sidecar 直接杀掉 Pod。
方案 3:共享进程命名空间(最终采用)
Use shared process namespace. This allows all pids in a pod to share the namespace. Thus pids can see each other.
利用 Pod 的shareProcessNamespace让 sidecar 与主容器共享 PID 空间。README 记录了两个待解问题及解法:
- 如何识别主容器:容器 ID 无法提前得知,容器名到 PID 的映射不可行;一种思路是调用 Kube API 获取 Pod 信息找到容器 ID,另一种更轻量的思路是轮询
/proc/<pid>/cgroup文件——该文件包含容器 ID,可以建立"盲"的容器 ID → PID 映射; - 如何等待主容器启动:同样可以借助 Kube API 获取容器 ID 后结合 cgroup 映射解决;
- 一旦拿到主容器,等待其退出、复制数据这两步"已实现"。
落地实现:Watcher 抽象
当前仓库把"等待启动/等待退出"抽象为containerwatcher.Watcher接口(flytecopilot/cmd/containerwatcher/iface.go):
type Watcher interface { WaitToStart(ctx context.Context) error WaitToExit(ctx context.Context) error }内置两种实现(--start-watcher-type/--exit-watcher-type可选):
- signal(默认):依赖 Kubernetes 1.28 的 sidecar container 原生特性,sidecar 监听
SIGTERM,收到信号即认为主容器退出(flytecopilot/cmd/containerwatcher/signal_watcher.go)。WaitToStart对 signal watcher 而言是 no-op; - noop:调试用占位实现,直接假定容器已启动/已退出(flytecopilot/cmd/containerwatcher/noop_watcher.go)。
协议层支撑:IOStrategy 与 DataLoadingConfig
CoPilot 的"何时下载/何时上传"策略由 flyteidl2/core/tasks.proto 中的IOStrategy定义,这也是 downloader/sidecar 各参数取值的权威来源:
| 枚举 | 取值 | 语义 |
|---|---|---|
DownloadMode | DOWNLOAD_EAGER(默认) | 主容器启动前全部下载完成 |
DownloadMode | DOWNLOAD_STREAM | 流式下载,写 End-Of-Stream 标记表示全部就绪 |
DownloadMode | DO_NOT_DOWNLOAD | 大对象(offloaded)不下载 |
UploadMode | UPLOAD_ON_EXIT(默认) | 主容器退出后统一上传 |
UploadMode | UPLOAD_EAGER | 数据出现即上传 |
UploadMode | DO_NOT_UPLOAD | 只写引用不上传数据 |
DataLoadingConfig(flyteidl2/core/tasks.proto)则定义了数据加载的完整开关与路径约定:
enabled:总开关,未设置则不启用数据加载;input_path/output_path:输入下载目录 / 输出上传目录(从根开始的绝对路径);format:JSON/YAML/PROTO三种元数据编码——有 protobuf 定义时推荐PROTO,无定义时JSON/YAML更易读;io_strategy:下载/上传时机策略;file_input_layout:DIRECT(默认,File 直接落到input_path/<var>,list[File] 落到input_path/<var>/<index>,丢弃原始文件名与扩展名)/NAMED_DIR(每个 File 放入独立目录并保留扩展名,list[File] 元素为input_path/<var>/<index><ext>,便于依赖扩展名识别格式的工具按 glob 消费)。
NAMED_DIR布局在下载器中的实现可参见 flytecopilot/data/download.go(list[File] 保留原始 basename,重名时回退为索引前缀)与RecursiveDownload中单 File 的目录化处理(同文件 L525-L531)。
数据层实现亮点
data包是 CoPilot 的数据搬运核心(flytecopilot/data/common.go 的注释将其定位为"目前只有两个工具:downloader 与 uploader"),几个值得关注的实现细节:
- 多部分 Blob 并发下载:
handleBlob对MULTIPART类型先通过store.List分页(每批 100 项)递归列出全部 part,再以 goroutine + WaitGroup 并发下载,并用 Mutex 保护目录创建与计数统计;任何文件下载失败或读写流未正常关闭都会汇总为明确错误(flytecopilot/data/download.go); - HTTP 直下载:对
http/https协议的引用走DownloadFileFromHTTP(带 context 取消的 GET 请求),其余走DownloadFileFromStorage(先 Head 确认存在再 ReadRaw),兼容 S3/GCS 等 stow 后端(flytecopilot/data/utils.go); - Offloaded 大对象处理:
RecursiveDownload在下载每个变量前检查literal.GetOffloadedMetadata(),若存在则先从远程读取真正的字面量内容再继续下载(flytecopilot/data/download.go); - 简单类型按原生格式落盘:
handlePrimitive支持 string、bool、integer、float、datetime(RFC3339Nano)、duration 等原生类型序列化写入(flytecopilot/data/download.go)。
错误处理闭环
CoPilot 的错误处理设计保证了"主容器失败也能被平台感知":
- 主容器写入
_ERROR文件,内容可为结构化core.ErrorDocument或纯文本; - Uploader 上传前先检查该文件:结构化错误 →
RawContainerError,由 sidecar 原样上传错误文档(kind/origin 由调用方设置,见 flytecopilot/cmd/root.go),平台据此判断是否可重试; - 非结构化错误 → 包装为
User Error: <内容>上传; - downloader/sidecar 自身的失败(如下载失败)则通过
UploadError(ctx, "InputDownloadFailed" / "OutputUploadFailed", ...)写入error.pb(错误文件名由根命令的--err-output-name控制,默认error.pb)。
限制与注意事项
- 简单类型输出文件上限 1024 字节,超出会报错(flytecopilot/data/upload.go);
- Uploader 目前仅支持
LiteralType_Blob与LiteralType_Simple两类输出类型,其余类型返回 "currently CoPilot uploader does not support ... system error"(flytecopilot/data/upload.go); - 重试尚未在
data层自动处理(见 flytecopilot/data/common.go 的 TODO 注释); DataLoadingConfig的注释明确指出该能力仅支持 Kubernetes("This is supported only on K8s at the moment")。
源码阅读地图
若想深入理解,推荐按以下路径阅读:
- 协议定义:flyteidl2/core/tasks.proto(IOStrategy / DataLoadingConfig / FileInputLayout)
- CLI 入口与根命令:flytecopilot/main.go、flytecopilot/cmd/root.go
- 下载器命令与实现:flytecopilot/cmd/download.go、flytecopilot/data/download.go
- 上传器命令与实现:flytecopilot/cmd/sidecar.go、flytecopilot/data/upload.go
- 容器生命周期 watcher:flytecopilot/cmd/containerwatcher/iface.go、flytecopilot/cmd/containerwatcher/signal_watcher.go
总结而言,Flyte CoPilot 的优雅之处在于:它不要求容器理解任何 Flyte 专有协议,只要求容器读写本地文件系统——下载器把远程输入翻译成本地文件,sidecar 把本地输出翻译回远程元数据,中间的主容器则完全保持"原生"。
- 后端
- 任务调度
- 工作流自动化
- 云原生
- MLOps
- 微服务
【免费下载链接】flyte
Dynamic, resilient AI orchestration. Coordinate data, models, and compute as you build AI workflows.
相关推荐
Flyte项目原生调度器架构深度解析
Flyte项目原生调度器架构深度解析 概述:云原生工作流编排的核心引擎 Flyte Propeller(螺旋桨)是Flyte项目的核心调度引擎,作为Kubern
后端任务调度工作流自动化云原生MLOps微服务Distrobox 实战指南:在任意 Linux 主机上运行任意发行版容器
Distrobox 实战指南:在任意 Linux 主机上运行任意发行版容器 Distrobox 是一个基于 podman 、 docker 或 lilipod
开发工具CLI探索downkyicore:用专业工具轻松提取B站音频的实战指南
探索downkyicore:用专业工具轻松提取B站音频的实战指南 downkyicore作为一款功能强大的哔哩哔哩视频下载工具,不仅支持8K、HDR、杜比视界等
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考