使用 AWS S3 设置自动化模型训练工作流
工作流自动化的开源方法
https://khuyentran1476.medium.com/?source=post_page---byline--cd0587b42f34--------------------------------https://towardsdatascience.com/?source=post_page---byline--cd0587b42f34-------------------------------- Khuyen Tran
·发布于 Towards Data Science ·阅读时间 7 分钟 ·2024 年 3 月 18 日
–
动机
假设你是一个电子商务平台,旨在增强推荐个性化。你的数据存储在 S3 中。
为了优化推荐,你计划每当 S3 中添加新文件时,使用最新的客户互动数据重新训练推荐模型。那么,如何具体处理这个任务呢?
除非另有说明,所有图片均由作者提供
解决方案
解决这个问题的两种常见方法是:
-
AWS Lambda: AWS 提供的无服务器计算服务,允许在事件触发时执行代码,无需管理服务器。
-
开源调度器: 自动化、调度和监控工作流与任务的工具,通常是自托管的。
使用开源调度器相较于 AWS Lambda 提供了以下优势:
-
成本效益: 在 AWS Lambda 上运行长时间任务可能会很昂贵。开源调度器允许你使用自己的基础设施,潜在地节省成本。
-
更快的迭代: 在本地开发和测试工作流能加快过程,使调试和优化变得更加容易。
-
环境控制: 对执行环境的完全控制使你能够根据自己的喜好自定义开发工具和集成开发环境(IDE)。
虽然你可以在 Apache Airflow 中解决这个问题,但它需要复杂的基础设施和部署设置。因此,我们将使用 Kestra,它提供了直观的用户界面,并且可以通过一个 Docker 命令启动。
欢迎在此处播放并分叉本文的源代码:
[## GitHub - khuyentran1401/mlops-kestra-workflow
通过在 GitHub 上创建账户,参与 khuyentran1401/mlops-kestra-workflow 的开发。
github.com](https://github.com/khuyentran1401/mlops-kestra-workflow?source=post_page-----cd0587b42f34--------------------------------)
工作流摘要
该工作流由两个主要组件组成:Python 脚本和编排。
编排
- Python 脚本和流程存储在 Git 中,并且会按计划同步到 Kestra。
- 当 S3 桶的“new”前缀下出现新文件时,Kestra 会触发一系列 Python 脚本的执行。
Python 脚本
-
download_files_from_s3.py: 从桶中的“old”前缀下载所有文件。
-
merge_data.py: 合并下载的文件。
-
process.py: 处理合并的数据。
-
train.py: 使用处理后的数据训练模型。
由于我们将在 Kestra 中执行从 Git 下载的代码,请确保将这些 Python 脚本提交到仓库中。
git add .
git commit -m 'add python scripts'
git push origin main
编排
启动 Kestra
执行以下命令下载 Docker Compose 文件:
curl -o docker-compose.yml \
https://raw.githubusercontent.com/kestra-io/kestra/develop/docker-compose.yml
确保 Docker 正在运行。然后,使用以下命令启动 Kestra 服务器:
docker-compose up -d
通过在浏览器中打开 URL localhost:8080 访问 UI。
从 Git 同步
由于 Python 脚本托管在 GitHub 上,我们将使用 Git Sync 每分钟将代码从 GitHub 同步到 Kestra。要设置此功能,请在“_flows”目录下创建一个名为“sync_from_git.yml”的文件。
.
├── _flows/
│ └── sync_from_git.yml
└── src/
├── download_files_from_s3.py
├── helpers.py
├── merge_data.py
├── process.py
└── train.py
如果您使用的是 VSCode,可以使用 Kestra 插件 来启用 .yaml 文件中的流程自动完成和验证功能。
https://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/5d963b494550c1f0572c20a04a2f7675.pnghttps://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/0d1942b43ca808d76b4744901d3bb420.png
以下是从 Git 同步代码的流程实现:
id: sync_from_git
namespace: dev
tasks:
- id: git
type: io.kestra.plugin.git.Sync
url: https://github.com/khuyentran1401/mlops-kestra-workflow
branch: main
username: "{{secret('GITHUB_USERNAME')}}"
password: "{{secret('GITHUB_PASSWORD')}}"
dryRun: false # if true, you'll see what files will be added, modified
# or deleted based on the Git version without overwriting the files yet
triggers:
- id: schedule
type: io.kestra.core.models.triggers.types.Schedule
cron: "*/1 * * * *" # every minute
只有在 GitHub 仓库是私有的情况下,才需要提供用户名和密码。要将这些密钥传递给 Kestra,请将它们放入“.env”文件中:
# .env
GITHUB_USERNAME=mygithubusername
GITHUB_PASSWORD=mygithubtoken
AWS_ACCESS_KEY_ID=myawsaccesskey
AWS_SECRET_ACCESS_KEY=myawssecretaccesskey
# ! This line should be empty
接下来,使用以下 bash 脚本对这些密钥进行编码:
while IFS='=' read -r key value; do
echo "SECRET_$key=$(echo -n "$value" | base64)";
done < .env > .env_encoded
执行此脚本会生成一个包含编码后的密钥的“.env_encoded”文件:
# .env_encoded
SECRET_GITHUB_USERNAME=bXlnaXRodWJ1c2VybmFtZQ==
SECRET_GITHUB_PASSWORD=bXlnaXRodWJ0b2tlbg==
SECRET_AWS_ACCESS_KEY_ID=bXlhd3NhY2Nlc3NrZXk=
SECRET_AWS_SECRET_ACCESS_KEY=bXlhd3NzZWNyZXRhY2Nlc3NrZXk=
在 Docker Compose 文件中包含编码后的环境文件,以便 Kestra 访问环境变量:
# docker-compose.yml
kestra:
image: kestra/kestra:latest-full
env_file:
- .env_encoded
确保在“.gitignore”文件中排除环境文件:
# .gitignore
.env
.env_encoded
最后,将新流程和 Docker Compose 文件都提交到 Git:
git add _flows/sync_from_git.yml docker-compose.yml
git commit -m 'add Git Sync'
git push origin main
现在,随着sync_from_git流程设置为每分钟运行一次,你可以方便地通过 Kestra UI 访问并触发 Python 脚本的执行。
https://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/720f5af3075edd1455253930ff0f021c.pnghttps://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/84846040a61cc99923ac42d464e5516c.png
编排
我们将创建一个流程,当有新文件添加到“winequality-red”桶中的“new”前缀时触发。
一旦检测到新文件,Kestra 会将其下载到内部存储并执行 Python 文件。最后,它将文件从“new”前缀移动到“old”前缀,以避免在随后的轮询中重复检测。
id: run_ml_pipeline
namespace: dev
tasks:
- id: run_python_commands
type: io.kestra.plugin.scripts.python.Commands
namespaceFiles:
enabled: true
env:
AWS_ACCESS_KEY_ID: "{{secret('AWS_ACCESS_KEY_ID')}}"
AWS_SECRET_ACCESS_KEY: "{{secret('AWS_SECRET_ACCESS_KEY')}}"
docker:
image: ghcr.io/kestra-io/pydata:latest
beforeCommands:
- pip install -r requirements.txt
commands:
- python src/download_files_from_s3.py
- python src/merge_data.py
- python src/process.py
- python src/train.py
outputFiles:
- "*.pkl"
triggers:
- id: watch
type: io.kestra.plugin.aws.s3.Trigger
interval: PT1S
accessKeyId: "{{secret('AWS_ACCESS_KEY_ID')}}"
secretKeyId: "{{secret('AWS_SECRET_ACCESS_KEY')}}"
region: us-east-2
bucket: winequality-red
prefix: new
action: MOVE
moveTo:
bucket: winequality-red
key: old
run_python_commands任务使用:
-
使用
namespaceFiles访问本地项目中的所有文件,并与 Git 仓库同步。 -
使用
env来检索环境变量。 -
使用
docker在 docker 容器ghcr.io/kestra-io/pydata:latest内执行脚本。 -
使用
beforeCommands在执行命令之前从“requirements.txt”文件安装依赖。 -
使用
commands按顺序执行命令列表。 -
使用
outputFiles将所有 pickle 文件从本地文件系统发送到 Kestra 的内部存储。
最后,添加upload任务,将模型的 pickle 文件上传到 S3。
id: run_ml_pipeline
namespace: dev
tasks:
- id: run_python_commands
type: io.kestra.plugin.scripts.python.Commands
namespaceFiles:
enabled: true
env:
AWS_ACCESS_KEY_ID: "{{secret('AWS_ACCESS_KEY_ID')}}"
AWS_SECRET_ACCESS_KEY: "{{secret('AWS_SECRET_ACCESS_KEY')}}"
docker:
image: ghcr.io/kestra-io/pydata:latest
beforeCommands:
- pip install -r requirements.txt
commands:
- python src/download_files_from_s3.py
- python src/merge_data.py
- python src/process.py
- python src/train.py model_path=model/model.pkl
outputFiles:
- "*.pkl"
# ------------------------- ADD THIS ------------------------- #
- id: upload
type: io.kestra.plugin.aws.s3.Upload
accessKeyId: "{{secret('AWS_ACCESS_KEY_ID')}}"
secretKeyId: "{{secret('AWS_SECRET_ACCESS_KEY')}}"
region: us-east-2
from: '{{outputs.run_python_commands.outputFiles["model/model.pkl"]}}'
bucket: winequality-red
key: model.pkl
# ------------------------------------------------------------ #
triggers:
...
就这样!将此流程命名为“run_ml_pipeline.yml”并提交到 Git。
git add run_ml_pipeline.yml
git commit -m 'add run_ml_pipeline'
git push origin main
触发流程
要启动流程,只需将新文件添加到 S3 的“winequality-red”桶中的“new”前缀。
该操作将触发run_ml_pipeline流程,启动从“旧”前缀下载数据、合并所有文件、处理数据并训练模型。
https://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/6e9a0b6ab67bc14bc15e7eca149714ec.pnghttps://github.com/OpenDocCN/towardsdatascience-blog-zh-2024/raw/master/docs/img/d3cfe0e9fa37c79f185715567c6f35f9.png
一旦工作流执行完毕,“model.pkl”文件会被上传到 S3。
结论
本文展示了如何使用 Kestra 自动化执行数据科学任务的 Python 脚本,每当有新文件添加到 S3 时。如果你在寻找自动化机器学习管道的方式,可以尝试这个解决方案。
我喜欢写关于数据科学的概念,并玩弄各种数据科学工具。你可以通过以下方式保持关注我的最新文章:
-
订阅我的Data Science Simplified新闻通讯。
AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。
更多推荐



所有评论(0)