Skip to content
mumong's blog
Go back

Kubeflow 安装部署与使用

记录 Kubeflow 在 Kubernetes 集群中的安装部署、存储准备、Notebook、Pipeline、AutoML、多租户和模型服务化使用流程。

本文由内部 DOCX 技术文档整理而来,已移除目录前的模板页与前置信息;正文内容、截图与操作步骤按博客阅读方式重新排版。

Kubeflow 是一个用于在 Kubernetes 上部署、管理和扩展机器学习(ML)工作流的开源平台。它的设计目标是简化 ML 项目在生产环境中的部署和操作,通过本文将介绍如何安装部署使用kubeflow

本文档适用于在一个初始化的k8s集群中部署使用kubeflow

术语、缩略语解释

引用文件

https://developer.aliyun.com/article/856853

https://github.com/kubeflow/manifests/tree/v1.8.1

https://developer.aliyun.com/article/784044

https://www.51cto.com/article/614445.html

https://www.cnblogs.com/muzinan110/p/18031944

https://www.cnblogs.com/tongai/p/17759117.html

https://www.wyx.cloudns.asia/blog/2022/01/28/kubeflow

https://docs.nvidia.com/datacenter/cloud-native/gpu-operator/latest/getting-started.html

什么是kubeflow

Kubeflow 是一个专为在 Kubernetes 上部署、管理和扩展机器学习(ML)工作流而设计的开源平台,旨在简化机器学习项目在生产环境中的部署和操作。基于Kubernetes的容器编排和资源管理功能。通过将机器学习工作流拆分为一系列的容器化任务,Kubeflow可以利用Kubernetes的自动扩展、容错和调度功能,确保机器学习任务的高效执行,Kubeflow 提供了一套完整的工具和服务,支持从数据准备、模型训练、调优到部署的整个机器学习生命周期。这种端到端的解决方案能够帮助数据科学家和工程师快速实现从开发到生产的转换。

Kubeflow是的机器学习工具包。Kubeflow是运行在K8S之上的一套技术栈,这套技术栈包含了很多组件,组件之间的关系比较松散,我们可以配合起来用,也可以单独用其中的一部分。下图是官网显示Kubeflow作为在上安排ML系统组件的平台:

Kubeflow 安装部署与使用 图 1
图 1

开发ML系统是一个反复的过程。我们需要评估ML工作流各个阶段的输出,并在必要时对模型和参数进行更改,以确保模型不断产生所需的结果。

Kubeflow 安装部署与使用 图 2
图 2
Kubeflow 安装部署与使用 图 3
图 3

由此可以看出,Kubeflow的目标是基于K8S,构建一整套统一的机器学习平台,覆盖最主要的机器学习流程(数据->特征->建模->服务→监控),同时兼顾机器学习的实验探索阶段和正式的生产环境。

可以直接从Kubeflow中的Jupyter notebook创建诸如deployment,Tfjob,service,pod等资源。notebook中已预装了命令行工具,可以说也是非常简单了。 将Jupyter notebook绑定在Kubeflow中时,可以使用Fairing库使用TFJob提交训练作业。训练作业可以运行在单个节点,也可以分布在同一个集群上,但不能在notebook pod内部运行。通过Fairing库提交作业可以使数据科学家清楚地了解Docker容器化和pod分配等流程。总体而言,Kubeflow-hosted notebooks可以更好地与其他组件集成,同时提供notebook image的可扩展性。

Kubeflow的目的主要是为了简化在上运行机器学习任务的流程,最终希望能够实现一套完整可用的流水线, 来实现机器学习从数据到模型的一整套端到端的过程。 而pipeline是一个工作流平台,能够编译部署机器学习的工作流。所以从这个层面来说,pipeline能够成为Kubeflow的核心组件一点也不意外。 kubeflow/pipelines实现了一个工作流模型。所谓工作流,或者称之为流水线(pipeline),其中的每一个节点被称作组件(component)。组件处理真正的逻辑,比如预处理,数据清洗,模型训练等。每一个组件负责的功能不同,但有一个共同点,即组件都是以Docker镜像的方式被打包,以容器的方式被运行的。 下图显示了Kubeflow Pipelines UI中管道的运行时执行图:

Kubeflow 安装部署与使用 图 4
图 4

流水线的定义可以分为两步,首先是定义组件,组件可以从镜像开始完全自定义。这里介绍一下自定义的方式:首先需要打包一个Docker镜像,这个镜像是组件的依赖,每一个组件的运行,就是一个Docker容器。其次需要为其定义一个python函数,描述组件的输入输出等信息,这一定义是为了能够让流水线理解组件在流水线中的结构,有几个输入节点,几个输出节点等。接下来组件的使用就与普通的组件并无二致了。 实现流水线的第二步,就是根据定义好的组件组成流水线,在流水线中,由输入输出关系会确定图上的边以及方向。在定义好流水线后,可以通过 python中实现好的流水线客户端提交到系统中运行。

Kubeflow提供了Katib组件,方便用户在Kubernetes集群上轻松执行超参优化。Katib的灵感来自Google的黑盒优化框架Vizier。它利用贝叶斯优化等先进的搜索算法来寻找最优的超参配置

Kubeflow 安装部署与使用 图 5
图 5

通过在jupyter中或者直接编辑yaml文件来定义和创建experiment,通过设置期望的优化算法来实现对超参数的调优与最优化结果的解析。

整体来说kubeflow有以下组件构成

下面给出部署kubeflow的一些环境准备工作。

Kubeflow 安装部署与使用 图 6
图 6

在部署使用kubeflow之前,我们需要有一个可以使用的k8s集群,在本示例中使用的k8s集群为1.26.8版本,对应kubeflow版本为1.8.1

Kubeflow 安装部署与使用 图 7
图 7

Nfs文件系统准备

由于安装部署kubeflow的要求是需要有一个default的storage class所以现在给出如何安装部署使用nfs文件系统与nfs-subdir-external-provisioner。

Kubeflow 安装部署与使用 图 8
图 8

在官方文档中,使用kubeflow mainfest安装的需求是有默认的storageclass,kustomize,与kubectl工具。

由于要使用默认的storage class因此部署使用nfs-subdir-external-provisioner。可以自动的创建和管理pv,pvc关系。因此介绍下对应的nfs,nfs插件的安装部署流程。

Nfs安装

通过命令apt install -y nfs-kernel-server安装下载nfs服务所以集群节点中。

创建一个数据存储目录 mkdir -p /data/redis。

修改/etc/exports文件下的内容,添加新的配置目录

Kubeflow 安装部署与使用 图 9
图 9

/data/redis 192.168.0.0/24(rw,sync,no_all_squash,no_subtree_check,no_root_squash)

完成后重启nfs的服务,使其处于running状态。

systemctl restart nfs-kernel-server.service

部署NFS-Subdir-External-Provisioner

NFS-Subdir-External-Provisioner是一个自动配置卷程序,它使用现有的和已配置的 NFS 服务器来支持通过持久卷声明动态配置 Kubernetes 持久卷。

在部署使用NFS-subdir服务前你需要一个有nfs服务的节点,并且知道对应的存储目录,在上述示例中他所对应的节点和目录为192.168.0.208,/data/redis

创建serviceaccount

展开代码片段(29 行)
apiVersion: v1
kind: ServiceAccount
metadata:
name: nfs-client-provisioner
namespace: default # 替换成你要部署的 Namespace
---
kind: ClusterRole
apiVersion: rbac.authorization.k8s.io/v1
metadata:
name: nfs-client-provisioner-runner
rules:
- apiGroups: [""]
resources: ["persistentvolumes"]
verbs: ["get", "list", "watch", "create", "delete"]
- apiGroups: [""]
resources: ["persistentvolumeclaims"]
verbs: ["get", "list", "watch", "update"]
- apiGroups: ["storage.k8s.io"]
resources: ["storageclasses"]
verbs: ["get", "list", "watch"]
- apiGroups: [""]
resources: ["events"]
verbs: ["create", "update", "patch"]
---
kind: ClusterRoleBinding
apiVersion: rbac.authorization.k8s.io/v1
metadata:
name: run-nfs-client-provisioner
subjects:

- kind: ServiceAccount

name: nfs-client-provisioner
namespace: default
roleRef:
kind: ClusterRole
name: nfs-client-provisioner-runner

apiGroup: rbac.authorization.k8s.io

展开代码片段(17 行)
---
kind: Role
apiVersion: rbac.authorization.k8s.io/v1
metadata:
name: leader-locking-nfs-client-provisioner
namespace: default
rules:
- apiGroups: [""]
resources: ["endpoints"]
verbs: ["get", "list", "watch", "create", "update", "patch"]
---
kind: RoleBinding
apiVersion: rbac.authorization.k8s.io/v1
metadata:
name: leader-locking-nfs-client-provisioner
namespace: default
subjects:

- kind: ServiceAccount

name: nfs-client-provisioner
namespace: default
roleRef:
kind: Role
name: leader-locking-nfs-client-provisioner

apiGroup: rbac.authorization.k8s.io

部署 NFS-Subdir-External-Provisioner,根据实际情况修改配置参数

apiVersion: apps/v1
kind: Deployment
metadata:
name: nfs-client-provisioner
labels:

app: nfs-client-provisioner

spec:

replicas: 1

strategy:

type: Recreate ## 设置升级策略为删除再创建(默认为滚动更新)

selector:
matchLabels:

app: nfs-client-provisioner

template:
metadata:
labels:

app: nfs-client-provisioner

展开代码片段(18 行)
spec:
serviceAccountName: nfs-client-provisioner
containers:
- name: nfs-client-provisioner
#image: gcr.io/k8s-staging-sig-storage/nfs-subdir-external-provisioner:v4.0.0
image: registry.cn-beijing.aliyuncs.com/xngczl/nfs-subdir-external-provisione:v4.0.0
volumeMounts:
- name: nfs-client-root
mountPath: /persistentvolumes
env:
- name: PROVISIONER_NAME ## Provisioner的名称,以后设置的storageclass要和这个保持一致
value: nfs-client
- name: NFS_SERVER ## NFS服务器地址,需和valumes参数中配置的保持一致
value: 192.168.0.208
- name: NFS_PATH ## NFS服务器数据存储目录,需和valumes参数中配置的保持一致
value: /data/redis
volumes:
- name: nfs-client-root

nfs:

server: 192.168.0.208 ## NFS服务器地址
path: /data/redis ## NFS服务器数据存储目录

创建 NFS StorageClass

apiVersion: storage.k8s.io/v1
kind: StorageClass
metadata:
name: nfs-storage
annotations:

storageclass.kubernetes.io/is-default-class: "true" ## 是否设置为默认的storageclass

provisioner: nfs-client ## 动态卷分配者名称,必须和上面创建的"provisioner"变量中设置的Name一致

parameters:

archiveOnDelete: "true" ## 设置为"false"时删除PVC不会保留数据,"true"则保留数据

之后应用上述的三个文件,等待拉取镜像并且部署使用

执行命令sudo chmod 777 /data/redis/ 为nfs存储的目录设置权限。

安装kubeflow

在官方gitbub仓库的kubeflow manifest项目中下载对应的二进制文件安装包skustomize_v5.0.3_linux_amd64.tar.gz,解压到/usr/local/bin目录下。

sudo tar -xzf kustomize_v5.0.3_linux_amd64.tar.gz -C /usr/local/bin

下载官方kubeflow maiinfest项目包mainfests-1.8.1.tar.gz

Kubeflow 安装部署与使用 图 10
图 10

修改默认存储sc

修改yaml,下面每个文件里面添加 storageClassName: nfs-storage,在mainfests-1.8.1目录下的下面四个文件

Kubeflow 安装部署与使用 图 11
图 11

apps/katib/upstream/components/mysql/pvc.yaml

common/oidc-client/oidc-authservice/base/pvc.yaml

apps/pipeline/upstream/third-party/minio/base/minio-pvc.yaml

apps/pipeline/upstream/third-party/mysql/base/mysql-pv-claim.yaml

修改镜像拉取策略

由于部分组件的镜像拉取策略为Always,所以修改他们为IfNotPresent,在当前目录下执行命令。find ./ -type f -exec grep -l "imagePullPolicy: Always" {} \;

查找到所有关于imagePullPolicy为Always的文件。在这一步可以不修改,在后续部署pod的过程中如果有pod所在节点存在镜像但是还是无法拉取镜像并且运行的话,就修改对应服务的镜像拉取策略。一键执行命令如下:可以全部替换

展开代码片段(20 行)
find ./ -type f -exec sed -i 's/imagePullPolicy: Always/imagePullPolicy: IfNotPresent/g' {} \;
root@master:~/huhu/kubeflow/manifests-1.8.1# find ./ -type f -exec grep -l "imagePullPolicy: Always" {} \;
./docs/KustomizeBestPractices.md
./contrib/kserve/models-web-app/base/deployment.yaml
./contrib/kserve/kserve/kserve.yaml
./contrib/kserve/kserve/kserve_kubeflow.yaml
./apps/kfp-tekton/upstream/third-party/tekton-custom-task/driver-controller/500-controller.yaml
./apps/kfp-tekton/upstream/third-party/kfp-csi-s3/csi-s3-deployment.yaml
./apps/kfp-tekton/upstream/v1/third-party/kfp-csi-s3/csi-s3-deployment.yaml
./apps/kfp-tekton/upstream/v1/base/cache-deployer/cache-deployer-deployment.yaml
./apps/kfp-tekton/upstream/v1/base/cache/cache-deployment.yaml
./apps/kfp-tekton/upstream/v1/base/pipeline/ml-pipeline-apiserver-deployment.yaml
./apps/kfp-tekton/upstream/v1/base/pipeline/ml-pipeline-viewer-crd-deployment.yaml
./apps/kfp-tekton/upstream/base/cache-deployer/cache-deployer-deployment.yaml
./apps/kfp-tekton/upstream/base/cache/cache-deployment.yaml
./apps/kfp-tekton/upstream/base/pipeline/ml-pipeline-viewer-crd-deployment.yaml
./apps/pipeline/upstream/base/cache-deployer/cache-deployer-deployment.yaml
./apps/pipeline/upstream/base/cache/cache-deployment.yaml
./apps/pipeline/upstream/base/pipeline/ml-pipeline-viewer-crd-deployment.yaml
./common/oidc-client/oidc-authservice/base/statefulset.yaml
Kubeflow 安装部署与使用 图 12
图 12

修改APP_SECURE_COOKIES

修改对应配置文件将APP_SECURE_COOKIES的值设置为false表示不使用加密cookies交互。一键执行命令如下:可以全部替换find ./ -type f -exec sed -i '/SECURE_COOKIES/s/=true/=false/g' {} \;

Vim ./apps/jupyter/jupyter-web-app/upstream/base/params.env
Kubeflow 安装部署与使用 图 13
图 13

将其修改为false,否则部署起来后无法通过dashbord访问kubeflow。

同样可以执行find ./ -type f -exec grep -l "APP_SECURE_COOKIES" {} \;命令来查看是否还有其他需要修改的APP_SECURE_COOKIES=false的配置,在我的部署过程中目前看来只需要修改jwp的即可。

root@master:~/huhu/kubeflow/manifests-1.8.1# find ./ -type f -exec grep -l "APP_SECURE_COOKIES" {} \;
./README.md
./apps/jupyter/jupyter-web-app/upstream/base/deployment.yaml
./apps/jupyter/jupyter-web-app/upstream/base/kustomization.yaml
./apps/jupyter/jupyter-web-app/upstream/base/params.env
./apps/volumes-web-app/upstream/base/deployment.yaml
./apps/volumes-web-app/upstream/base/kustomization.yaml
./apps/volumes-web-app/upstream/base/params.env
./apps/tensorboard/tensorboards-web-app/upstream/base/deployment.yaml
./apps/tensorboard/tensorboards-web-app/upstream/base/kustomization.yaml
./apps/tensorboard/tensorboards-web-app/upstream/base/params.env

部署kubeflow

在修改完上述对应的配置后,在mainfest目录下执行命令如下可以自动的部署kubeflow,由于它所依赖的组件过多,因此安装过程中大概率会出现问题,上述的配置应该能解决大部分遇到的问题,但是如果遇到安装失败,pod无法正常Running请自行排查解决。

执行命令

while ! kustomize build example | kubectl apply -f -; do echo "Retrying to apply resources"; sleep 10; done

等待pod全部running,期间会出现错误,逐步排查解决对应的pod问题,大多都是镜像无法拉取,可以手动拉取并且分发上传到对应节点上。同时看需求修改对应deployment,statusfulset的镜像拉取策略为IfNotPresent,就可以修复异常pod。

镜像拉取失败

由于国内使用网络环境的问题,在拉取大部分的镜像的时候可能都会出现失败的情况,因此需要自行解决镜像拉取的问题,如可以使用代理的方式手动下载所需的镜像。下面给出一个可以使用的脚本用于将master节点所缺少的镜像同步下载分发到不同节点。

前提是需要sshpass工具。并且节点配置了代理。该工具通过查看所有pod不是running状态并且所需的镜像名字,通过脚本的方式对所需镜像进行下载和打包转发到其他node节点。在脚本中需要配置workernode和用户登录的账户密码。

但是由于脚本的不成熟性,对于某些镜像可能还需要手动下载并上传到worker节点使用。

#!/bin/bash

# 定义工作节点IP
WORKER_NODES=("")
USER="root"
PASS="huhu"
# 检查sshpass是否安装
if ! command -v sshpass &> /dev/null
then

echo "sshpass could not be found. Please install it first."

exit 1

fi

# 显示进度条函数
show_progress() {
local duration=$1
local sleep_interval=0.1
local progress=0
local bar_size=40
while [ $progress -lt 100 ]; do
local num_bars=$((progress * bar_size / 100))
printf "\r[%-${bar_size}s] %d%%" $(printf "#%.0s" $(seq 1 $num_bars)) $progress
progress=$((progress + 1))

sleep $sleep_interval

done

echo

}
# 补全镜像路径函数
get_full_image_path() {
local image=$1
# 如果包含 / 或 .io,直接返回
if [[ "$image" == *"/"* || "$image" == *".io"* ]]; then

echo "$image"

else

echo "docker.io/library/$image"

fi

# 处理 kserve/models-web-app 这种情况
if [[ "$image" != *".io"* && "$image" != *"/"* ]]; then

echo "docker.io/$image"

fi

}
echo "==============="

echo "检查需要下载的镜像"

echo "==============="

declare -A all_images

# 获取异常状态的容器镜像
get_problem_images() {
local namespace=$1
local pod_name=$2
# 获取所有容器的镜像(包括init容器和普通容器)

local all_container_images=$(kubectl -n $namespace get pod $pod_name -o jsonpath='{range .spec.initContainers[*]}{.image}{"\n"}{end}{range .spec.containers[*]}{.image}{"\n"}{end}')

echo "$all_container_images"

}
# 获取不完全ready的pods
while IFS= read -r line; do
if [[ -n "$line" ]]; then
namespace=$(echo "$line" | awk '{print $1}')
pod_name=$(echo "$line" | awk '{print $2}')
status=$(echo "$line" | awk '{print $4}')
# 检查各种异常状态
if [[ "$status" == "ImagePullBackOff" ||
"$status" == "ErrImagePull" ||
"$status" == "CreateContainerError" ||
"$status" == "CrashLoopBackOff" ||
"$status" == "Error" ||
"$status" == "ContainerCreating" ]]; then

echo "检查异常状态的Pod: $namespace/$pod_name (Status: $status)"

while IFS= read -r image; do
if [[ -n "$image" ]]; then
if [[ "$image" != *:* ]]; then
image="${image}:latest"
fi
full_image_path=$(get_full_image_path "$image")
if [[ "$full_image_path" != *".io"* && "$full_image_path" != "docker.io"* ]]; then
full_image_path="docker.io/$full_image_path"
fi
all_images["$full_image_path"]=1

echo "发现问题镜像: $full_image_path"

fi

done < <(get_problem_images "$namespace" "$pod_name")
fi
fi
done < <(kubectl get pods -A | grep -v "Running\|Completed\|NAMESPACE")
corrected_images="${!all_images[@]}"
if [ -z "$corrected_images" ]; then

echo "没有找到需要下载的镜像。"

exit 0

fi

echo "列出所需的镜像:"

for image in $corrected_images; do

echo "- $image"

done

echo

echo "==============="

echo "下载所需镜像"

echo "==============="
for image in $corrected_images; do
safe_filename=$(echo "$image" | sed 's/[\/:]/-/g').tar
if [ -f "/tmp/$safe_filename" ]; then

echo "跳过下载镜像 $image, 因为它已经存在于/tmp目录中"

continue

fi

echo "正在拉取镜像 $image"

if ctr -n k8s.io i pull "$image"; then
ctr -n k8s.io i export "/tmp/$safe_filename" "$image"

show_progress 2

else

echo "拉取镜像 $image 失败,尝试替代镜像源..."

# 尝试 docker.io 镜像
alternative_image="docker.io/${image}"

echo "尝试拉取: $alternative_image"

if ctr -n k8s.io i pull "$alternative_image"; then
ctr -n k8s.io i export "/tmp/$safe_filename" "$alternative_image"

show_progress 2

else
# 尝试 gcr.io 镜像
alternative_image="gcr.io/${image#*/}"

echo "尝试拉取: $alternative_image"

if ctr -n k8s.io i pull "$alternative_image"; then
ctr -n k8s.io i export "/tmp/$safe_filename" "$alternative_image"

show_progress 2

fi

fi

fi

done

for i in "${!WORKER_NODES[@]}"; do
worker=${WORKER_NODES[$i]}
node_num=$((i+1))
echo "==========="
echo "传输到 node${node_num} (${worker})"
echo "==========="
for image in $corrected_images; do
safe_filename=$(echo "$image" | sed 's/[\/:]/-/g').tar
if [ -f "/tmp/$safe_filename" ]; then
sshpass -p ${PASS} scp -o StrictHostKeyChecking=no "/tmp/$safe_filename" "${USER}@${worker}:/tmp/"
sshpass -p ${PASS} ssh -o StrictHostKeyChecking=no ${USER}@${worker} "ctr -n k8s.io i import '/tmp/$safe_filename' && rm '/tmp/$safe_filename'"
fi
done
done
echo "==========="

echo "清理临时文件"

echo "==========="
for image in $corrected_images; do
safe_filename=$(echo "$image" | sed 's/[\/:]/-/g').tar

rm "/tmp/$safe_filename"

echo "已删除 /tmp/$safe_filename"

done

修改svc为nodeport模式

由于部署起来的svc全都是cluser ip的类型所以无法直接被外部访问,因此需要手动修改istio-ingressgateway的svc为nodeport类型

k edit svc -n istio-system istio-ingressgateway

Kubeflow 安装部署与使用 图 14
图 14

之后查看nodeport对应的端口,通过IP:端口的方式可以在页面端访问kubeflow

Kubeflow 安装部署与使用 图 15
图 15

访问kubeflow 对应的IP+30981 手动修改后会随机的分配一个端口

默认的登录账户密码为

user@example.com

12341234

Kubeflow 安装部署与使用 图 16
图 16

使用kubeflow

Kubeflow 安装部署与使用 图 17
图 17

这里按照模块介绍下 Kubeflow 的几个核心组件。

Kubeflow 安装部署与使用 图 18
图 18
Kubeflow 安装部署与使用 图 19
图 19

使用的过程中确保pod的服务都是running 的状态,由于kubeflow安装有很多个组件,因此可以根据自己的需求添加其他不同的组件,在官方github项目中有的对应的project的安装部署方式。

在部署使用一个服务后,kubeflow会自动的在kubeflow-user-example-com命名空间下创建任务pod如nfs-test-0.

Kubeflow 安装部署与使用 图 20
图 20
Kubeflow 安装部署与使用 图 21
图 21

在他的内部可以正常使用功能,如示例中使用python运行了简单的代码和功能。

模型开发-Notebooks

Kubeflow Notebooks 是 Kubeflow 平台中一个关键的组件,它为数据科学家和机器学习工程师提供了在 Kubernetes 上运行 Jupyter Notebooks 的能力。这一功能的出现极大地简化了在云环境中管理和使用 Jupyter Notebook 的复杂性。通过 Kubeflow Notebooks,用户能够在 Kubernetes 集群中轻松创建和管理多个 Jupyter Notebook 实例。这些实例可以针对特定用户进行资源配置,如 CPU、内存和 GPU,以确保在多用户环境中能够实现高效的资源隔离和使用。

Kubeflow 安装部署与使用 图 22
图 22

通过notebook的模块新建新的task

Kubeflow 安装部署与使用 图 23
图 23

在配置页面可以设置名你在,对应的cpu,gpu的调配,创建好之后会出现对应的pod。关于gpu的使用可以参考官方文档的介绍https://v1-8-branch.kubeflow.org/docs/components/notebooks/quickstart-guide/

Kubeflow 安装部署与使用 图 24
图 24
Kubeflow 安装部署与使用 图 25
图 25

点击connect链接到对应的pod中,在这一步会遇到镜像拉取失败的问题,需要手动拉取镜像到指定节点等待pod运行起来。

ctr -n k8s.io i pull xxx
Kubeflow 安装部署与使用 图 26
图 26

JupyterLab 提供了一个高度集成的工作环境,支持各种工具和文件类型,从而帮助用户执行从数据科学到软件开发的各种任务。其中,Jupyter Notebook 是该平台的核心功能之一,允许用户创建包含可执行代码、富文本、图表以及其他动态元素的文档。这些笔记本支持多种编程语言,并广泛应用于数据分析、数学建模和统计可视化等领域。

除了笔记本,JupyterLab 还包括一个交互式控制台,用于快速执行代码段,非常适合进行临时实验或调试。此外,它还提供了一个集成的终端,使用户能够直接访问操作系统的命令行功能,进行文件管理或系统配置等操作。这一功能尤其对那些需要直接与底层系统交互的用户来说非常有用。

对于需要编辑纯文本或脚本的用户,JupyterLab 提供了文本文件编辑器和专门的 Python 文件编辑器,这使得用户能够在同一环境中开发软件或编写代码。Markdown 文件编辑器则允许用户创建和编辑轻量级的标记文档,非常适合编写具有丰富格式的文档和报告。此外,"Show Contextual Help" 工具能够在侧边栏显示活动代码的相关文档,帮助用户快速理解正在使用的代码库或API的细节。

模型训练

Kubeflow 安装部署与使用 图 27
图 27

在创建notebook的时候可以进行镜像的选择,在这里我们选择带有tensorflow的镜像,就可以直接在里面使用对应的框架。同时还有不同的如pytorch-cuda镜像等提供。创建好后通过如下的示例来跑一个简单的训练代码。

import numpy as np
import tensorflow as tf
from tensorflow.keras.datasets import mnist
from tensorflow.keras.utils import to_categorical
from tensorflow.keras.models import Sequential
from tensorflow.keras.layers import Dense, Flatten, Dropout, Conv2D, MaxPooling2D
# 加载MNIST数据集
(train_images, train_labels), (test_images, test_labels) = mnist.load_data()
# 预处理数据:调整形状并归一化
train_images = train_images.reshape(-1, 28, 28, 1).astype('float32') / 255.0
test_images = test_images.reshape(-1, 28, 28, 1).astype('float32') / 255.0
# 将标签转换为 one-hot 编码
train_labels = to_categorical(train_labels, 10)
test_labels = to_categorical(test_labels, 10)
# 构建模型
model = Sequential([
Conv2D(32, kernel_size=(3, 3), activation='relu', input_shape=(28, 28, 1)),
MaxPooling2D(pool_size=(2, 2)),
Conv2D(64, kernel_size=(3, 3), activation='relu'),
MaxPooling2D(pool_size=(2, 2)),
Flatten(),
Dense(128, activation='relu'),
Dropout(0.5),
Dense(10, activation='softmax')
])
# 编译模型
model.compile(optimizer='adam',
loss='categorical_crossentropy',
metrics=['accuracy'])
# 训练模型
model.fit(train_images, train_labels, epochs=10, batch_size=64, validation_data=(test_images, test_labels))
# 评估模型
test_loss, test_accuracy = model.evaluate(test_images, test_labels, verbose=0)
print(f"Test accuracy: {test_accuracy:.4f}")
# 保存模型
model.save('mnist_cnn_model.keras')
print("模型保存成功")

在jupyter中运行可以得到对应的结果,这样我们就完成一个简单的模型训练效果。

Kubeflow 安装部署与使用 图 28
图 28

训练完成后可以同时的使用训练后的模型来进行预测。Model.save()方式保存模型。在新的jupyter文件中通过import模型来导入。使用自己手写的图片来进行结果的预测。

import numpy as np
import tensorflow as tf
from PIL import Image, ImageOps, ImageEnhance, ImageFilter
import matplotlib.pyplot as plt
# 加载保存的模型
loaded_model = tf.keras.models.load_model('mnist_cnn_model.keras')
# 定义你自己的图片文件名列表
image_files = ['9.jpg', '7.jpg','3.jpg', 'third.jpg']
def preprocess_image(img):
# 反转颜色:白底黑字 -> 黑底白字
img = ImageOps.invert(img)
# 转换为灰度图像
img = img.convert('L')
# 自动对比度增强
img = ImageOps.autocontrast(img)
# 应用二值化
img = img.point(lambda p: 255 if p > 128 else 0)
# 使用 ImageFilter 模拟膨胀效果(类似于最大滤波)
img = img.filter(ImageFilter.MaxFilter(3))
# 使用 ImageFilter 模糊处理
img = img.filter(ImageFilter.GaussianBlur(1))
# 再次膨胀处理
img = img.filter(ImageFilter.MaxFilter(3))
# 再次进行对比度增强
img = ImageEnhance.Contrast(img).enhance(2.0)
# 锐化处理
img = img.filter(ImageFilter.SHARPEN)
# 裁剪数字的边缘并居中
img = img.crop(img.getbbox()) # 裁剪非空白区域
img = img.resize((20, 20), Image.Resampling.LANCZOS) # 调整图像大小,保持最大信息
background = Image.new('L', (28, 28), 0) # 创建黑色背景
offset = ((28 - img.size[0]) // 2, (28 - img.size[1]) // 2)
background.paste(img, offset) # 将图像粘贴到背景上使其居中
return background
# 创建一个图形,包含4行2列的子图
fig, axs = plt.subplots(4, 2, figsize=(10, 20))
for i, file in enumerate(image_files):
# 加载图像
img = Image.open(file)
# 对图像进行预处理
processed_img = preprocess_image(img)
# 将图像转换为数组并进行标准化,确保形状为 (28, 28, 1)
img_array = np.array(processed_img).reshape(1, 28, 28, 1).astype('float32') / 255.0
# 进行预测
predictions = loaded_model.predict(img_array)
# 获取预测结果
predicted_digit = np.argmax(predictions[0])
# 显示原始图片
axs[i, 0].imshow(img, cmap='gray')
axs[i, 0].set_title(f'原始图片: {file}')
axs[i, 0].axis('off')
# 显示处理后的图片
axs[i, 1].imshow(processed_img, cmap='gray')
axs[i, 1].set_title(f'预测结果: {predicted_digit}')
axs[i, 1].axis('off')
print(f"{file} 预测的数字是: {predicted_digit}")
plt.tight_layout()
plt.show()
Kubeflow 安装部署与使用 图 29
图 29
Kubeflow 安装部署与使用 图 30
图 30

不过由于模型比较小,目前看起来预测的准确率不太行,后续可以继续优化使用更好的模型和训练数据。

同时也可以使用mnist里面的图片进行预测,这样结果的准确率就会高很多。

Kubeflow 安装部署与使用 图 31
图 31

GPU训练

使用gpu镜像会要求我们的集群中存在GPU资源如下所示

Kubeflow 安装部署与使用 图 32
图 32

当我们选中gpu镜像,想要添加gpu使会提示集群中不存在gpu,所以需要在某个节点插上物理gpu然后再集群中添加operator来使用gpu资源

下面给出安装nvidia-gpu operator的方法:

参考官方安装nvidia-operator链接

下载准备helm3

curl -fsSL -o get_helm.sh https://raw.githubusercontent.com/helm/helm/master/scripts/get-helm-3 \
&& chmod 700 get_helm.sh \
&& ./get_helm.sh

确保NFD模式是关闭的,如果有开启的那么手动关闭它

kubectl get nodes -o json | jq '.items[].metadata.labels | keys | any(startswith("feature.node.kubernetes.io"))'

添加helm仓库

helm repo add nvidia https://helm.ngc.nvidia.com/nvidia \
&& helm repo update

部署gpu-operator

helm install --wait --generate-name \
-n gpu-operator --create-namespace \
nvidia/gpu-operator \
--set driver.version=535

其中需要注意的是对应的驱动版本我用的是A800因此他是535.根据自己的nvidia-gpu型号确定自己的驱动版本。

上面的流程执行完成后说明operator已经安装完成。

Kubeflow 安装部署与使用 图 33
图 33

拥有gpu的node就会出现新的可调度资源nvidia.com/gpu

Kubeflow 安装部署与使用 图 34
图 34

打开kubelfow开始使用gpu训练任务。需要再这个页面指定含有cuda的镜像,并且在gpu配置中选择集群中可用的gpu nvidia。

Kubeflow 安装部署与使用 图 35
图 35

创建好后进入jupyter内创建一个python3工具。在其中可以进行机器学习代码的开发。同时可以使用nvidia-smi命令在pod内部查看和适用到我们的GPU

Kubeflow 安装部署与使用 图 36
图 36

一切准备就绪后我们就可以开始用jupyter进行模型训练

下面给出训练的代码

import os
import tensorflow as tf
import time
import numpy as np
from tensorflow.keras import layers, models
import matplotlib.pyplot as plt
print("====检查 GPU 可用性====")
os.environ['TF_FORCE_GPU_ALLOW_GROWTH'] = 'true'
# 检查 GPU 是否可用
if tf.test.is_gpu_available():
print("\033[1;32m[GPU 可用] 将进行 GPU 和 CPU 训练对比\033[0m")
gpu_device = tf.config.list_physical_devices('GPU')[0]
print(f"可用的 GPU: {gpu_device}")
else:
print("\033[1;31m[GPU 不可用] 只能使用 CPU 进行训练\033[0m")
exit()
print("\n====加载和预处理数据====")
(train_images, train_labels), (test_images, test_labels) = tf.keras.datasets.mnist.load_data()
train_images = train_images.reshape((60000, 28, 28, 1)).astype('float32') / 255
test_images = test_images.reshape((10000, 28, 28, 1)).astype('float32') / 255
展开代码片段(15 行)
def create_model():
model = models.Sequential([
layers.Conv2D(32, (3, 3), activation='relu', input_shape=(28, 28, 1)),
layers.MaxPooling2D((2, 2)),
layers.Conv2D(64, (3, 3), activation='relu'),
layers.MaxPooling2D((2, 2)),
layers.Conv2D(64, (3, 3), activation='relu'),
layers.Flatten(),
layers.Dense(64, activation='relu'),
layers.Dense(10, activation='softmax')
])
model.compile(optimizer='adam',
loss='sparse_categorical_crossentropy',
metrics=['accuracy'])
return model
# GPU 训练
print("\n====开始 GPU 训练====")
with tf.device('/GPU:0'):
gpu_model = create_model()
start_time = time.time()
gpu_history = gpu_model.fit(train_images, train_labels, epochs=10,
validation_split=0.2, batch_size=64, verbose=1)
gpu_time = time.time() - start_time
# CPU 训练
print("\n====开始 CPU 训练====")
os.environ['CUDA_VISIBLE_DEVICES'] = '-1' # 禁用 GPU
with tf.device('/CPU:0'):
cpu_model = create_model()
start_time = time.time()
cpu_history = cpu_model.fit(train_images, train_labels, epochs=10,
validation_split=0.2, batch_size=64, verbose=1)
cpu_time = time.time() - start_time
# 结果对比
print("\n====训练时间对比====")
print(f"\033[1;34mGPU 训练时间: {gpu_time:.2f} 秒\033[0m")
print(f"\033[1;34mCPU 训练时间: {cpu_time:.2f} 秒\033[0m")
print(f"\033[1;32mGPU 加速比: {cpu_time / gpu_time:.2f}x\033[0m")
# 绘制训练过程的损失和准确率曲线
def plot_history(history, title):
acc = history.history['accuracy']
val_acc = history.history['val_accuracy']
loss = history.history['loss']
val_loss = history.history['val_loss']
epochs = range(1, len(acc) + 1)
plt.figure(figsize=(12, 5))
plt.subplot(1, 2, 1)
plt.plot(epochs, loss, 'bo-', label='Training loss')
plt.plot(epochs, val_loss, 'ro-', label='Validation loss')
plt.title(f'{title} - Training and validation loss')
plt.xlabel('Epochs')
plt.ylabel('Loss')
plt.legend()
plt.subplot(1, 2, 2)
plt.plot(epochs, acc, 'bo-', label='Training accuracy')
plt.plot(epochs, val_acc, 'ro-', label='Validation accuracy')
plt.title(f'{title} - Training and validation accuracy')
plt.xlabel('Epochs')
plt.ylabel('Accuracy')
plt.legend()
plt.show()
print("\n====可视化 GPU 训练过程====")
plot_history(gpu_history, "GPU Training")
print("\n====可视化 CPU 训练过程====")
plot_history(cpu_history, "CPU Training")
# 评估 GPU 模型
print("\n====评估 GPU 训练的模型====")
test_loss, test_acc = gpu_model.evaluate(test_images, test_labels, verbose=0)
print(f'\n\033[1;32mTest accuracy: {test_acc:.4f}\033[0m')
# 保存 GPU 训练的模型
print("\n====保存 GPU 训练的模型====")
gpu_model.save('mnist_model_gpu.keras')
print("模型已保存为 mnist_model_gpu.keras")

这段代码的主要功能是通过对比 GPU 和 CPU 在相同任务上的训练表现,来体现 GPU 的强大计算能力,尤其是在深度学习任务中的显著优势。代码首先检查当前环境中是否可用 GPU,并根据设备的可用性分别在 GPU 和 CPU 上训练一个简单的卷积神经网络(CNN),该网络用于对 MNIST 手写数字数据集进行分类。

具体来说,代码加载并预处理了 MNIST 数据集,然后定义了一个包含多个卷积层、池化层和全连接层的卷积神经网络模型。在确认 GPU 可用的情况下,代码在 GPU 上训练模型,并记录训练时间。接着,禁用 GPU,只使用 CPU 进行相同的训练任务,并再次记录时间。通过比较 GPU 和 CPU 的训练时间,代码展示了 GPU 在处理深度学习任务时的加速效果。

此外,代码还通过绘制训练过程中的损失和准确率曲线,直观地展示了模型在训练和验证集上的表现。这不仅有助于理解模型的训练过程,还可以用于分析模型是否出现过拟合或欠拟合的现象。

在一个配备1个CPU和4GB内存的环境下运行这段代码,可以清楚地看到 GPU 在深度学习任务中的优势。即便在有限的资源配置下,GPU 仍能显著加速训练过程,从而提高了计算效率和任务完成的速度。通过这种对比分析,有效地展示了 GPU 的强大功能和在深度学习领域中的重要性。

重要部分为通过使用with tf.device('/GPU:0'): 来使用指定的GPU进行训练

Kubeflow 安装部署与使用 图 37
图 37

这是在同样规格下的训练见过

Kubeflow 安装部署与使用 图 38
图 38

同时在训练的过程中使用nvidia-smi命令可以查看到硬件GPU的使用情况

Kubeflow 安装部署与使用 图 39
图 39

比如上图中可以看出在训练过程中GPU的利用率分别在9%左右,说明我们的任务成功的调用了GPU进行计算。

Pipline

Kubeflow Pipelines 是 Kubeflow 项目中的一个核心模块,专注于构建、部署和管理复杂的机器学习工作流。它提供了一整套用于设计和自动化机器学习流水线的工具,使数据科学家和工程师能够更加高效地构建、管理和监控机器学习模型的训练和部署过程。通过 Kubeflow Pipelines,用户可以轻松地定义、分享、重用和自动化复杂的工作流。

在pipline中可以上传制作好的.gz文件yaml文件等。创建好流程图后会有如下图形界面上显示。具体的制作过程可以参考官方文档https://www.kubeflow.org/docs/components/pipelines/getting-started/

Kubeflow 安装部署与使用 图 40
图 40
import kfp
from kfp import dsl
from kfp.dsl import component, Input, Output, Dataset, Model
展开代码片段(19 行)
# Step 1: 数据下载和预处理
@component(
base_image='python:3.8-slim',
packages_to_install=[
'pandas',
'scikit-learn',
'joblib',
'numpy',
'requests'
]
)
def preprocess_data_op(output_data: Output[Dataset]):
print("开始执行 preprocess_data_op...")
try:
import pandas as pd
print("成功导入 pandas 模块。")
except ImportError as e:
print(f"导入 pandas 失败: {e}")
raise e
from sklearn.model_selection import train_test_split
from sklearn.preprocessing import StandardScaler
import os
print("正在下载数据集...")
url = "https://raw.githubusercontent.com/jbrownlee/Datasets/master/pima-indians-diabetes.data.csv"
columns = ['Pregnancies', 'Glucose', 'BloodPressure', 'SkinThickness', 'Insulin', 'BMI', 'DiabetesPedigreeFunction', 'Age', 'Outcome']
data = pd.read_csv(url, names=columns)
print("数据集下载完成。")
# 数据清洗和特征工程
print("正在进行数据清洗和特征工程...")
X = data.drop('Outcome', axis=1)
y = data['Outcome']
# 标准化特征
print("正在标准化特征...")
scaler = StandardScaler()
X_scaled = scaler.fit_transform(X)
# 划分训练集和测试集
print("正在划分训练集和测试集...")
X_train, X_test, y_train, y_test = train_test_split(X_scaled, y, test_size=0.2, random_state=42)
# 保存预处理后的数据
print("正在保存预处理后的数据...")
os.makedirs(output_data.path, exist_ok=True)
pd.DataFrame(X_train).to_csv(os.path.join(output_data.path, 'X_train.csv'), index=False)
pd.DataFrame(X_test).to_csv(os.path.join(output_data.path, 'X_test.csv'), index=False)
pd.DataFrame(y_train).to_csv(os.path.join(output_data.path, 'y_train.csv'), index=False)
pd.DataFrame(y_test).to_csv(os.path.join(output_data.path, 'y_test.csv'), index=False)
print(f"数据预处理完成并已保存到 {output_data.path}。")
展开代码片段(18 行)
# Step 2: 模型训练
@component(
base_image='python:3.8-slim',
packages_to_install=[
'pandas',
'scikit-learn',
'joblib',
'numpy'
]
)
def train_model_op(input_data: Input[Dataset], output_model: Output[Model]):
print("开始执行 train_model_op...")
try:
import pandas as pd
print("成功导入 pandas 模块。")
except ImportError as e:
print(f"导入 pandas 失败: {e}")
raise e
from sklearn.linear_model import LogisticRegression
import joblib
import os
print("正在加载训练数据...")
X_train = pd.read_csv(os.path.join(input_data.path, 'X_train.csv'))
y_train = pd.read_csv(os.path.join(input_data.path, 'y_train.csv'))
# 训练模型
print("正在训练模型...")
model = LogisticRegression()
model.fit(X_train, y_train.values.ravel())
# 创建输出目录并保存模型
os.makedirs(output_model.path, exist_ok=True) # 确保输出目录存在
model_path = os.path.join(output_model.path, 'trained_model.joblib')
joblib.dump(model, model_path)
print(f"模型训练完成并已保存到 {model_path}。")
展开代码片段(18 行)
# Step 3: 模型评估
@component(
base_image='python:3.8-slim',
packages_to_install=[
'pandas',
'scikit-learn',
'joblib',
'numpy'
]
)
def evaluate_model_op(input_data: Input[Dataset], input_model: Input[Model]):
print("开始执行 evaluate_model_op...")
try:
import pandas as pd
print("成功导入 pandas 模块。")
except ImportError as e:
print(f"导入 pandas 失败: {e}")
raise e
from sklearn.metrics import accuracy_score
import joblib
import os # 添加os模块的导入
print("正在加载测试数据和模型...")
X_test = pd.read_csv(os.path.join(input_data.path, 'X_test.csv'))
y_test = pd.read_csv(os.path.join(input_data.path, 'y_test.csv'))
model = joblib.load(os.path.join(input_model.path, 'trained_model.joblib'))
# 预测和评估
print("正在进行模型预测和评估...")
y_pred = model.predict(X_test)
accuracy = accuracy_score(y_test, y_pred)
print(f"模型准确率: {accuracy}")
# Step 4: Pipeline 定义
@dsl.pipeline(
name='Diabetes Classifier Pipeline',
description='A pipeline to train and evaluate a diabetes classifier model'
)
def diabetes_pipeline():
preprocess = preprocess_data_op()
train = train_model_op(input_data=preprocess.outputs['output_data'])
evaluate = evaluate_model_op(input_data=preprocess.outputs['output_data'], input_model=train.outputs['output_model'])
# Compile the pipeline
if __name__ == "__main__":
kfp.compiler.Compiler().compile(diabetes_pipeline, 'diabetes_pipeline.yaml')

还需要运行pvc文件

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
name: kubeflow-test-pv
namespace: kubeflow-user-example-com
spec:

accessModes:

- ReadWriteOnce

resources:
requests:
storage: 128Mi
Kubeflow 安装部署与使用 图 41
图 41
Kubeflow 安装部署与使用 图 42
图 42
Kubeflow 安装部署与使用 图 43
图 43

这个Python文件主要定义了一个基于Kubeflow Pipelines (KFP) 的数据处理与机器学习模型训练与评估的Pipeline,用于处理糖尿病分类任务。以下是该文件的功能和作用总结:

第一步:数据下载和预处理

文件中首先定义了 preprocess_data_op 组件,用于下载糖尿病数据集并进行数据预处理。具体操作包括数据下载、数据清洗、特征标准化以及训练集和测试集的划分。预处理后的数据被保存到指定的输出路径中,以便后续的模型训练和评估使用。

第二步:模型训练

接下来,文件定义了 train_model_op 组件,用于使用预处理后的训练数据来训练一个逻辑回归模型。训练好的模型会被保存到指定的输出路径中。通过使用 scikit-learn 库,组件能够快速实现模型的训练,并为后续的模型评估步骤提供一个已训练的模型文件。

第三步:模型评估

文件中还定义了 evaluate_model_op 组件,用于加载测试数据和训练好的模型,并使用测试数据对模型进行预测和评估。该组件计算并输出模型的准确率,作为模型性能的衡量标准。评估结果可以用于判断模型的有效性。

第四步:Pipeline 定义与编译

最后,文件定义了一个名为 diabetes_pipeline 的Pipeline,将上述三个组件串联在一起。通过 dsl.pipeline 装饰器,将数据预处理、模型训练和模型评估的步骤组合成一个完整的工作流。该Pipeline最终被编译为一个YAML文件(diabetes_pipeline.yaml),可以在Kubeflow环境中部署和运行。

总体而言,这个示例展示了 Argo Workflow 的基本用法,通过定义简单的任务和依赖关系,实现了在 Kubernetes 环境中的自动化工作流管理

假设现在你想写一个机器学习的 pipeline,大概抽象成几个步骤。

读取数据 -> 进行训练 -> 保存模型

以上的步骤可以通过pipline来实现,同时可以通过对应的流程输出来很好的管理模型数据之间的输入输出关系,方便操作与控制。

构建好后通过create run来运行对应的pipline实现我们的预期功能

Kubeflow 安装部署与使用 图 44
图 44

运行后可以在后台查看对应的pod是否正常启动,来执行任务。

一个大型的pipline可以兼顾多个方面,通过对自己工作流的理解,将模型与数据处理模块化,可以提高利用率。复用不同的模块并且

Kubeflow 安装部署与使用 图 45
图 45

对于需要多次执行,或者不同训练任务有共同的预处理项目。可以根据名字来使用之前处理好的缓存。提升训练效率

AutoML

AutoML 是机器学习比较热的领域,主要用来模型自动优化和超参数调整,这里其实是用的 Katib来实现的,一个基于k8s的 AutoML 项目,详细见https://github.com/kubeflow/katib。

Katib 主要提供了 超参数调整(Hyperparameter Tuning),早停法(Early Stopping)和神经网络架构搜索(Neural Architecture Search) Katib 支持多种超参数优化算法,包括随机搜索、贝叶斯优化和 HyperBand。这些算法可以帮助开发者在给定的参数空间中高效地搜索最佳超参数配置。通过定义实验,用户可以指定要调整的超参数及其范围、优化目标(如准确率或损失),以及使用的搜索算法。Katib 还支持提前停止策略,允许在训练过程中根据模型性能动态终止不佳的实验,从而节省计算资源。此外,Katib 不仅限于超参数调优,还提供了神经网络结构搜索(Neural Architecture Search, NAS)的功能。尽管这一功能仍在不断完善中,但它为用户提供了探索不同网络架构的可能性,进一步提升模型的表现。

Kubeflow 安装部署与使用 图 46
图 46

这些功能结合在一起,使得 Kubeflow 成为一个功能全面的机器学习平台,能够支持从数据处理到模型训练和部署的完整流程。

通过点击edit的方式来编辑yaml文件从而进行动态的调整参数。

也可以点击jupyter新建一个python文件来进行模型参数的优化。流程如下:

Kubeflow 安装部署与使用 图 47
图 47

新建文件notebooks如6.1节所示。创建好一个py文件后,在terminal中执行下面命令下载所需的库。

pip install kubeflow-katib

下载完成后进入编辑目录填入下方脚本

Kubeflow 安装部署与使用 图 48
图 48
import kubeflow.katib as katib
# Step 1. Create an objective function.
def objective(parameters):
# Import required packages.
import time
time.sleep(5)
# Calculate objective function.
result = 4 * int(parameters["a"]) - float(parameters["b"]) ** 2
# Katib parses metrics in this format: <metric-name>=<metric-value>.
print(f"result={result}")
# Step 2. Create HyperParameter search space.
parameters = {
"a": katib.search.int(min=10, max=20),
"b": katib.search.double(min=0.1, max=0.2)
}
# Step 3. Create Katib Experiment.
katib_client = katib.KatibClient()
name = "tune-experiment"
katib_client.tune(
name=name,
objective=objective,
parameters=parameters,
objective_metric_name="result",
max_trial_count=12
)
# Step 4. Get the best HyperParameters.
print(katib_client.get_optimal_hyperparameters(name))

运行后可以在autoML中查看参数优化文件和整体流程。

Kubeflow 安装部署与使用 图 49
图 49
Kubeflow 安装部署与使用 图 50
图 50

这段代码通过定义一个目标函数来自动进行超参数优化,使用随机搜索算法对该目标函数的形式是 result=4a - b^2,其中参数 a 是整数,范围为 10 到 20,参数 b 是浮点数,范围为 0.1 到 0.2。首先,代码定义了这个目标函数,然后创建了超参数的搜索空间。接着,使用Kubeflow Katib的 KatibClient 启动实验,并指定目标函数、参数范围和最大实验次数(12次)。最后,通过 katib_client.get_optimal_hyperparameters(name) 方法获取并输出最优的超参数组合,从而实现自动化的超参数调优,以优化指定的目标函数。

Kubeflow 安装部署与使用 图 51
图 51

a的取值范围为10-20,b的取值范围为0.1-0.2。在这个条件下使用随机搜索算法找到F的最大值。

Kubeflow 安装部署与使用 图 52
图 52

还可以根据自己的需求,修改目标函数和对应的优化算法。通过这种方式可以方便的为机器学习算法中调优参数优化提供很大的便利。

下面给出一个更加复杂的计算式与使用其他算法计算的示例:

Kubeflow 安装部署与使用 图 53
图 53
import kubeflow.katib as katib
# Step 1. Create an objective function.
def objective(parameters):
# Import required packages.
import time
import math
time.sleep(5)
# Calculate a more complex objective function.
a = int(parameters["a"])
b = float(parameters["b"])
# 复杂的目标函数示例:结合多项式、对数和三角函数的组合
result = (3 * a ** 2 + 2 * b - math.sin(a * b)) / (1 + math.log(b + 0.1)) + math.sqrt(abs(a - b))
# Katib parses metrics in this format: <metric-name>=<metric-value>.
print(f"result={result}")
# Step 2. Create HyperParameter search space.
parameters = {
"a": katib.search.int(min=10, max=20),
"b": katib.search.double(min=0.1, max=0.2)
}
展开代码片段(12 行)
# Step 3. Create Katib Experiment.
katib_client = katib.KatibClient()
name = "tune-experiment-for-kalibbase"
katib_client.tune(
name=name,
objective=objective,
parameters=parameters,
algorithm_name="bayesianoptimization",
objective_metric_name="result",
#objective_type="minimize", # 指定最小化目标函数
max_trial_count=12
)
# Step 4. Get the best HyperParameters.
print(katib_client.get_optimal_hyperparameters(name))

与上述例子对比起来修改的点不多,只有result表达式,并且在katib_client.tune中添加字段algorithm_name="bayesianoptimization",就可以以贝叶斯算法计算表达式的最大值。

Kubeflow 安装部署与使用 图 54
图 54

多租户使用kubeflow

Kubelfow支持使用profile的方式创建新用户。编辑一下文件创建用户。

apiVersion: kubeflow.org/v1beta1
kind: Profile
metadata:
name: test # 用户namespace
spec:
owner:
kind: User
name: huhu@test.com # 用户名

其中 kind表示用户的权限。创建好后执行命令kubectl apply -f new-profile.yaml

通过Kubectl get profile来查看创建状态

Kubeflow 安装部署与使用 图 55
图 55

之后由于kubeflow是通过configmap来管理用户信息的。所以需要修改dex组件的configmap

Kubectl edit cm -n auth dex
Kubeflow 安装部署与使用 图 56
图 56

标记的位置存储着默认的user信息。新加一下内容:

- email: huhu@test.com
hash: $2y$12$sgKXohry9ZCcZE/AXtxbgOWPD/elNJTslR6I1bpi.ordASFdg3Vmi
username: huhu

其中email为刚刚创建的用户的账户,hash为登录的密码但是以hash的方式展示。

需要apt install python3-pip 安装pip。

创建密码方式前提为pip install passlib bcrypt 之后执行一下的python代码

python3 -c 'from passlib.hash import bcrypt; import getpass; print(bcrypt.using(rounds=12, ident="2y").hash(getpass.getpass()))'
Kubeflow 安装部署与使用 图 57
图 57

执行命令后输入想要加密的登录密码,如huhu,之后会生成一段hash值,将其填入上述的cm配置中。保存退出后,重启pod就可以用新用户登录kubeflow

Kubeflow 安装部署与使用 图 58
图 58

不同的用户可以在不同的命名空间下操作和适用kubeflow。也可以在manage contributors来进行不同用户间的合作管理

Kubeflow 安装部署与使用 图 59
图 59

kserve

KServe是一个开源的云原生模型服务平台,旨在简化在Kubernetes上部署和运行机器学习模型的过程,支持多种机器学习框架、具备弹性扩容能力。KServe通过定义简单的YAML文件,提供声明式的API来部署模型,使得配置和管理模型服务变得更加容易。

Kubeflow 安装部署与使用 图 60
图 60

如何将我们训练后的模型部署起来,在前面的例子中我们训练了一个模型。叫做mnist_cnn_model.keras。现在下面给出如何将这个模型部署使用起来。

Kubeflow 安装部署与使用 图 61
图 61
Kubeflow 安装部署与使用 图 62
图 62

在控制端通过jupyter的模型名称找到对应的pod。之后describe查找到模型对应的存储pvc

Kubeflow 安装部署与使用 图 63
图 63

之后编辑yaml文件,其中需要注意的是namespace应当与所运行的pod一致。同时在storageUri中正确填写对应的pvc名字与模型的存储路径。

apiVersion: "serving.kserve.io/v1beta1"
kind: "InferenceService"
metadata:
name: "mnist-cnn-service"
namespace: kubeflow-user-example-com
spec:
predictor:
tensorflow:
storageUri: "pvc://model-show-workspace/1/mnist_cnn_model.keras"

在保存的模型中执行下面代码,生成可被识别的模型

import tensorflow as tf
# 加载 .keras 模型
model = tf.keras.models.load_model('/home/jovyan/1/mnist_cnn_model.keras')
# 将模型保存为 SavedModel 格式
model.save('/home/jovyan/1', save_format='tf')
Kubeflow 安装部署与使用 图 64
图 64

之后应用文件等待pod正常运行

Kubeflow 安装部署与使用 图 65
图 65
Kubeflow 安装部署与使用 图 66
图 66

得到部署起来的用于预测的predict服务。可以将刚刚训练好的模型对外提供服务,通过curl请求对应的svc得到response。

在dasbord中也可以查看到部署的模型

Kubeflow 安装部署与使用 图 67
图 67

将kubeflow与向量数据库整milvus合使用

对于大量模型训练来说,有一个好的向量数据库可以更加方便的处理数据的输入和输出,因此现在给出将kubeflow与milvus向量数据库结合起来使用。

部署使用milvus

详见我的另一篇技术文档”向量化数据库milvus”

在jupyter中使用k8s集群中的milvus

由于milvus是部署在k8s集群中的,因此在kubeflow中的jupyter中使用需要将milvus准备好。下面给出一个直接的示例

from pymilvus import connections, Collection, FieldSchema, CollectionSchema, DataType
import numpy as np
# 连接到 Milvus 实例
connections.connect(alias="default", host="192.168.0.208", port="31011")
print("已连接到 Milvus")
# 创建集合
fields = [
FieldSchema(name="id", dtype=DataType.INT64, is_primary=True),
FieldSchema(name="vector", dtype=DataType.FLOAT_VECTOR, dim=128)
]
schema = CollectionSchema(fields, description="用于演示的集合")
collection = Collection(name="demo_collection_2", schema=schema) # 使用新的集合名称
print(f"集合 {collection.name} 已创建")
# 插入向量数据
vectors = np.random.random((20, 128)).astype(np.float32) # 插入20个向量
collection.insert([list(range(20)), vectors])
print("已将数据插入集合")
# 为向量字段创建索引
index_params = {
"metric_type": "L2",
"index_type": "IVF_FLAT", # 选择适合的索引类型
"params": {"nlist": 128}
}
collection.create_index(field_name="vector", index_params=index_params)
print("索引已创建")
# 加载集合
collection.load()
print("集合已加载到内存")
# 查询向量数据
search_params = {"metric_type": "L2", "params": {"nprobe": 10}}
query_vectors = np.random.random((5, 128)).astype(np.float32) # 查询5个向量
search_results = collection.search(query_vectors, "vector", search_params, limit=3)
for i, result in enumerate(search_results):
print(f"查询向量 {i+1} 的前三个结果: {result}")
Kubeflow 安装部署与使用 图 68
图 68

在jupyter中新建一个python文件,将上述代码复制进去。处理好依赖相关的库运行。

得到结果如上所示。同时在milvus的可视化界面中也可以查询到我们的数据。

Kubeflow 安装部署与使用 图 69
图 69

这段代码展示了如何使用 Milvus 向量数据库来管理和查询高维向量数据。具体而言,代码首先连接到 Milvus 实例,然后创建一个名为 demo_collection_2 的集合,其中包含两个字段:一个用于存储整数 ID,另一个用于存储128维的浮点数向量。接着,代码生成并插入随机向量数据,为向量字段创建索引,并将集合加载到内存中。最后,代码对新的随机向量执行相似性搜索,返回最接近的结果。在 AI 和机器学习领域,这种向量搜索功能通常用于实现图像检索、文本搜索、推荐系统等应用,通过快速查找相似数据来支持复杂的预测和分类任务。

Milvus存储非结构化数据

由于milvus是向量数据库,因此下面给出一个可以将图片作为非结构数据存储到数据库中的方式

from pymilvus import connections, Collection, FieldSchema, CollectionSchema, DataType
import numpy as np
from tensorflow.keras.applications.resnet50 import ResNet50, preprocess_input
from tensorflow.keras.preprocessing import image
from tensorflow.keras.models import Model
from PIL import Image
import matplotlib.pyplot as plt
import time
# 连接到 Milvus 实例
connections.connect(alias="default", host="192.168.0.208", port="31011")
print("已连接到 Milvus")
# 加载预训练的ResNet50模型并去掉最后的分类层,用于特征提取
base_model = ResNet50(weights='imagenet')
model = Model(inputs=base_model.input, outputs=base_model.get_layer('avg_pool').output)
# 动态生成集合名称,避免冲突
collection_name = f"image_collection_{int(time.time())}"
# 定义集合
fields = [
FieldSchema(name="id", dtype=DataType.INT64, is_primary=True, auto_id=True), # 使用 auto_id 自动生成ID
FieldSchema(name="vector", dtype=DataType.FLOAT_VECTOR, dim=2048), # ResNet50的输出是2048维向量
FieldSchema(name="img_path", dtype=DataType.VARCHAR, max_length=255) # 存储图片路径
]
schema = CollectionSchema(fields, description="用于存储图片向量的集合")
collection = Collection(name=collection_name, schema=schema)
print(f"集合 {collection.name} 已创建")
# 图片列表
image_paths = ["3.jpg", "7.jpg", "9.jpg", "third.jpg"]
# 提取并存储图片特征向量和路径
all_features = []
all_paths = []
for img_path in image_paths:
# 加载图片并进行预处理
img = image.load_img(img_path, target_size=(224, 224))
x = image.img_to_array(img)
x = np.expand_dims(x, axis=0)
x = preprocess_input(x)
# 提取图片的特征向量
features = model.predict(x).flatten().tolist() # 转换为列表以插入 Milvus
# 收集数据
all_features.append(features)
all_paths.append(img_path)
# 将所有图片的特征向量和路径存储到 Milvus 中
entities = [
all_features, # 向量数据
all_paths # 路径数据
]
collection.insert(entities)
print(f"图片已存储")
# 为向量字段创建索引
index_params = {
"metric_type": "L2",
"index_type": "IVF_FLAT",
"params": {"nlist": 128}
}
collection.create_index(field_name="vector", index_params=index_params)
print("索引已创建")
# 加载集合
collection.load()
print("集合已加载到内存")
# 查询与某张图片最相似的图片
query_img_path = "third.jpg" # 要查询的图片
# 加载查询图片并提取特征向量
img = image.load_img(query_img_path, target_size=(224, 224))
x = image.img_to_array(img)
x = np.expand_dims(x, axis=0)
x = preprocess_input(x)
query_features = model.predict(x).flatten().tolist()
# 查询与这张图片最相似的图片(只返回一个结果)
search_params = {"metric_type": "L2", "params": {"nprobe": 10}}
search_results = collection.search([query_features], "vector", search_params, limit=1, output_fields=["img_path"])
展开代码片段(15 行)
# 读取查询结果并展示图片
for hits in search_results:
for hit in hits:
result_img_path = hit.entity.get('img_path')
if result_img_path:
try:
result_img = Image.open(result_img_path)
plt.imshow(result_img)
plt.title(f"检索到的图片: {result_img_path}")
plt.axis('off')
plt.show()
except Exception as e:
print(f"无法打开图片 {result_img_path}: {e}")
else:
print("未找到 'img_path' 字段")

运行后得到结果如下所示,我们通过上述的代码存储了4张图片进入milvus由于是非结构化向量存储,所以我们首先需要对图片进行向量化处理,代码中通过一个模型的最后一层输入,将一个图片分解为2048维的向量。通过将4张图片都分解后用一个2048维度的数组表示一张图片

最后通过查询的方式,将我们想要的照片从数据库中查询出来。通过这个例子很好的展示了如何使用向量化数据库milvus以及如何将一个图片进行向量化存储与查询的方式。

Kubeflow 安装部署与使用 图 70
图 70

总结

Kubeflow 安装部署与使用 图 71
图 71

一个简单的machine learning运行流程如上所示

整个流水线包括以下几部分:

基于上述功能描述我们其实可以基于 kubeflow 的 pipeline 和 kfserving 功能轻松实现一个简单的 MLOps 流水线发布流程。Kubeflow 是一个开源的机器学习平台,专为 Kubernetes 设计,旨在简化机器学习工作流的部署和管理。它将多种机器学习工具和框架整合到一个统一的生态系统中,提供了从数据准备到模型训练、优化和部署的全生命周期管理。

Kubeflow 的设计理念是提供一个全面、易用、可扩展的机器学习平台,利用 Kubernetes 的核心优势,如自动化部署、扩展和管理容器化应用程序。对于机器学习项目来说,Kubeflow 不仅提高了开发和部署的效率,还确保了解决方案的可移植性和可维护性。对于希望在 Kubernetes 上运行机器学习工作负载的团队而言,Kubeflow 提供了强大的工具和资源,使得机器学习的创新和实施更加便捷和高效。

Kubeflow的工作全都可以在jupyter中完成,如可以在jupyter中创建pipline,创建kalib等。这些通过调用api的方式都可以直接创建生成对应的workflow。让机器学习工程师在使用的过程中方便的管理自己的模型并且可以便利的进行参数调优以及通过搭建构建pipline的方式让其完成一系列流水线的操作,如数据清洗,批量操作。


Share this post on:

Previous Post
可观测性平台调研与实践
Next Post
vLLM 工具调研