ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

基于 Pachyderm 的 GPT-2 推文生成流水线:从数据采集、模型训练到文本生成的端到端实践

基于 Pachyderm 的 GPT-2 推文生成流水线:从数据采集、模型训练到文本生成的端到端实践 数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载导读本指南围绕 examples/ml/gpt-2/README.md 展开介绍如何在 Pachyderm 中搭建一条由三个环节组成的机器学习流水线推文抓取tweets→ 模型微调train→ 推文生成generate利用 OpenAI 的 GPT-2 文本生成模型借助 gpt-2-simple 封装库训练出会发推的模型。读完本文你将掌握如何在 Pachyderm 的 PFS版本化文件系统之上用管道规范pipeline spec串联 Python 脚本、配置 GPU 与内存资源、启用自动扩缩容并通过pachctl完成从创建数据仓库到查看最终生成结果的完整闭环。本示例假设你已经把 Pachyderm 部署并启动完毕教程重点全部放在流水线本身的构建上。Pachyderm 每个新的 minor 版本都会带来较大的架构变化因此该仓库中的示例按分支维护master 分支使用 Pachyderm 2.1.x 系列、2.0.x 分支对应 2.0.x 系列、1.13.x 分支对应 1.13.x 系列使用前请确认与你的版本匹配。一、整条流水线的 DAG 结构这条流水线在 DAG 顶端是一个名为queries的输入数据仓库里面存放我们想要执行的 Twitter 查询条件每条查询一个文件。整个处理流程共三步tweets读取queries仓库中的查询条件调用 twitterscraper 抓取推文并写入带 GPT-2 专用分隔符的文本train读取tweets仓库的抓取结果下载 GPT-2 预训练权重并用 gpt-2-simple 做微调产出模型文件generate读取train仓库中的微调模型调用generate_to_file生成一批模仿原用户风格的推文。每一步都是一个独立的 Pachyderm 管道前一步的输出 commit 会作为后一步的输入 commit数据沿 DAG 自动流动——这就是 Pachyderm数据驱动的核心工作方式只要有新数据提交进来管道就会自动触发。二、第一步推文抓取tweets 管道2.1 抓取脚本实现抓取逻辑位于 tweets.py#!/usr/local/bin/python3 import os import twitterscraper as t for query in os.listdir(/pfs/queries/): with open(os.path.join(/pfs/queries, query)) as f: for q in f: q q.strip() # clean whitespace with open(os.path.join(/pfs/out, query), w) as out: for tweet in t.query_tweets(q): out.write(|startoftext| ) out.write(tweet.text) out.write( |endoftext| )这段代码大部分是标准的 Pachyderm 管道写法关键在于理解两个挂载路径/pfs/queries输入仓库queries被挂载的位置os.listdir遍历到的每个文件就是一条查询条件/pfs/out输出目录管道脚本写入这里的文件会被自动收集为该管道的输出 commit并成为下游管道的输入。代码里值得注意的细节是每条推文前后被注入了两个特殊标记|startoftext|和|endoftext|。这两个标记是 GPT-2 在原始训练语料中就已经学习过的分隔符注入后可以让模型在生成阶段按一条推文为单位产出文本。仓库中实际的 tweets.py 在写入推文时还额外做了一次 ASCII 编码替换tweet.text.encode(ascii, replace).decode(ascii)用于过滤掉 Twitter 文本中的非 ASCII 字符避免对下游训练造成干扰。2.2 管道规范与部署命令对应的管道规范是 tweets.pipeline.json{ pipeline: { name: tweets }, description: A pipeline that scrapes tweets from https://twitter.com., transform: { image: pachyderm/gpt-2-example, cmd: [/tweets.py] }, input: { pfs: { repo: queries, glob: /* } } }关键配置说明pipeline.name管道名称也是下游读取该管道输出的仓库名transform.image管道运行所用的容器镜像此处为仓库自带的pachyderm/gpt-2-example构建方式见下文 Makefile 部分transform.cmd容器启动后执行的命令即抓取脚本input.pfs.repo输入仓库input.pfs.glob/*表示对输入仓库的每个一级目录/文件各启动一个 datum这样多条查询就能被并行处理。glob 模式定义在 src/pps/pps.proto 的PFSInput消息中是 PPS 将数据切分成处理单元datum的核心机制。在创建管道之前需要先创建输入仓库$ pachctl create repo queries然后创建管道$ pachctl create pipeline -f tweets.pipeline.json2.3 喂入数据并验证管道创建好后向queries仓库提交一条简单的 Twitter 查询即可触发任务。例如抓取某个用户名下的全部推文$ echo from:username | pachctl put file queriesmaster:username注意用户名不能包含符号。如果希望构造更复杂的查询可以使用 Twitter 的高级搜索功能点搜索按钮后页面顶部的查询字符串就是可用的q。提交之后queries仓库会产生一个新 committweets管道随即产生一个对应的输出 commit 和 job。查看运行中的任务$ pachctl list job任务结束后查看抓取到的推文$ pachctl get file tweetsmaster:/username确认结果合理后就可以进入模型训练环节。三、第二步模型微调train 管道3.1 训练脚本实现训练脚本位于 train.py使用 OpenAI 的 GPT-2 文本生成模型——准确说是它的便捷封装库 gpt-2-simple#!/usr/local/bin/python3 import gpt_2_simple as gpt2 import os tweets [f for f in os.listdir(/pfs/tweets)] # chdir so that the training process outputs to the right place out os.path.join(/pfs/out, tweets[0]) os.mkdir(out) os.chdir(out) model_name 117M gpt2.download_gpt2(model_namemodel_name) sess gpt2.start_tf_sess() gpt2.finetune(sess, os.path.join(/pfs/tweets, tweets[0]), model_namemodel_name, steps25) # steps is max number of training steps与第一步类似这里通过/pfs/tweets读取输入这次输入仓库是tweets。脚本做了两个关键选择模型规格默认使用 117M 参数的 GPT-2 版本。追求更好效果可以换成 345M 版本但训练时间会显著变长训练步数steps25这个取值在效果尚可与运行时间可接受之间做了折中属于示例性参数可按需调整。脚本中os.chdir(out)是 gpt-2-simple 的限制导致的它不支持指定模型输出目录只能通过切换当前工作目录的方式让训练产物checkpoint、模型文件落到/pfs/out下。这是 Pachyderm 管道脚本中很常见的一个技巧——当依赖库只认相对路径时用 chdir 把输出重定向到/pfs/out。需要说明的是README 中展示的这份脚本是教程形态的简化版117M、25 步仓库中实际的 train.py 已经演化为 345M 模型、1000 训练步两者结论一致仅参数不同。以你当前检出的代码为准。3.2 带资源限制与自动扩缩容的管道规范训练管道的规范是 train.pipeline.json{ pipeline: { name: train }, description: A pipeline that trains the ML model on the tweets gathered by the tweets pipeline., transform: { image: pachyderm/gpt-2-example, cmd: [/train.py] }, input: { pfs: { repo: tweets, glob: /* } }, resourceLimits: { gpu: { type: nvidia.com/gpu, number: 1 }, memory: 10G, cpu: 1 }, resourceRequests: { memory: 10G, cpu: 1 }, autoscaling: true }与tweets管道相比有三个变化输入换成tweets仓库transform 改为运行/train.py新增resourceLimits/resourceRequests模型微调是计算密集型任务需要显式申请 GPUnvidia.com/gpu数量 1、10G 内存和 1 核 CPU。这两类字段在 src/pps/pps.proto 的管道消息中分别对应ResourceSpec resource_requests与ResourceSpec resource_limits前者是调度的最低保障后者是运行时的硬性上限开启autoscaling: true防止管道在无数据处理时还占着 GPU 等资源不放空闲时自动释放、有数据时再拉起对按量付费的 GPU 集群尤为实用。autoscaling字段同样定义在 src/pps/pps.protobool autoscaling 33;。创建管道$ pachctl create pipeline -f train.pipeline.json由于tweets仓库中已经有待处理的数据管道创建后会立即触发一个 job。这个任务耗时较长README 作者在笔记本上运行约 1 小时想跑快一点可以调小steps或自行构建精简的 Docker 镜像。趁训练运行期间我们来搭建最后一步文本生成。四、第三步推文生成generate 管道4.1 生成脚本实现生成脚本位于 generate.py#!/usr/local/bin/python3 import gpt_2_simple as gpt2 import os models [f for f in os.listdir(/pfs/train)] model_dir os.path.join(/pfs/train, models[0]) # cant tell gpt2 where to read from, so we chdir os.chdir(model_dir) sess gpt2.start_tf_sess() gpt2.load_gpt2(sess) out os.path.join(/pfs/out, models[0]) gpt2.generate_to_file(sess, destination_pathout, prefix|startoftext|, truncate|endoftext|, include_prefixFalse, length280, nsamples30)除了读取/pfs/train训练管道的输出仓库、同样用os.chdir解决 gpt-2-simple 路径问题的样板代码外核心是generate_to_file这一调用它完成了真正的推文生成。几个参数的含义prefix|startoftext|提示模型从推文开头标记开始生成truncate|endoftext|在遇到推文结束标记时截断保证每次输出恰好是一条完整的推文——这正是第一步抓取时注入分隔符的原因模型把这些标记当作文本边界include_prefixFalse不让|startoftext|被重复追加到每一条生成结果里length280对应 Twitter 的单条推文长度上限 280 字符README 提到未来版本可能训练模型生成推文风暴 tweet storm即连发多条推文nsamples30生成 30 个样本本场景下每个样本就是一条推文。同样仓库中实际的 generate.py 参数已演进为nsamples200并增加了temperature1.0采样温度控制README 展示的 30 个样本是教程版参数两者机制一致。4.2 生成管道规范生成管道规范是 generate.pipeline.json结构上与 train 管道高度相似只是输入仓库换成了train、脚本换成了/generate.py{ pipeline: { name: generate }, description: A pipeline that generates tweets based on the trained model., transform: { image: pachyderm/gpt-2-example, cmd: [/generate.py] }, input: { pfs: { repo: train, glob: /* } }, resourceLimits: { gpu: { type: nvidia.com/gpu, number: 1 }, memory: 10G, cpu: 1 }, resourceRequests: { memory: 10G, cpu: 1 }, autoscaling: true }加载模型推理同样是 GPU 密集型工作因此保留了与 train 一致的资源申请和自动扩缩容配置。创建命令$ pachctl create pipeline -f generate.pipeline.json至此三条管道全部就绪形成queries → tweets → train → generate的完整数据流抓取的推文越多、越新训练出的模型就越贴合目标账号的说话风格生成的推文也就越逼真。五、修改与重新部署Makefile 与 Dockerfile示例自带一个简单的 Makefile 来构建与部署整套流水线docker-build: docker build -t pachyderm/gpt-2-example . docker-push: docker-build docker push pachyderm/gpt-2-example deploy: pachctl update repo queries pachctl create pipeline -f tweets.pipeline.json pachctl create pipeline -f train.pipeline.json pachctl create pipeline -f generate.pipeline.json修改代码后重新构建镜像$ make docker-build部署整条流水线make deploy会先update repo queries再依次创建三条管道$ make deploy如果修改了代码但希望复用既有镜像标签需要先make docker-build必要时docker-push推送镜像再用pachctl update pipeline让 Pachyderm 使用新镜像。镜像本身的构建方式是 DockerfileFROM tensorflow/tensorflow:1.14.0-gpu-py3 RUN apt-get update \ apt-get install -y python3-pip \ rm -rf /var/lib/apt/lists/* RUN pip3 install twitterscraper gpt_2_simple ADD tweets.py / ADD train.py / ADD generate.py /基础镜像选用 TensorFlow 1.14.0 的 GPU 版保证gpt_2_simple依赖 TensorFlow 1.x API开箱即用通过pip3安装twitterscraper与gpt_2_simple两个运行时依赖三个 Python 脚本被直接放进镜像根目录/与管道规范里的cmd: [/tweets.py]等命令路径一一对应。值得补充的是仓库中还有一个 README 未详述的 listen.py它以无限循环的方式持续监听某个话题示例为#PachydermODSC并把结果打包成 tar 归档写入/pfs/out可作为持续数据源型管道的参考实现但其并未出现在主流程的 Makefile 部署列表中。六、原理纵深这条示例背后的 Pachyderm 机制6.1 数据驱动触发与/pfs挂载整条流水线无需任何定时器或手动调度pachctl put file产生 commit → PPS 监听该 commit → 自动为下游管道创建 job。这种commit 即事件的模型是 Pachyderm 区别于传统批处理框架的核心。管道脚本统一从/pfs/repo读输入、向/pfs/out写输出PFS 层自动把输出收录为新的版本化 commit形成可追溯的数据血缘。6.2 glob 模式与 datum 并行三个管道的输入都使用glob: /*含义是按输入仓库的顶层条目切分 datum每个 datum 独立调度、可并行执行。glob 语义定义在 src/pps/pps.proto 的PFSInput消息中string glob 5;而Input联合消息src/pps/pps.proto决定了管道如何组合多路输入。对 tweets 管道而言/*意味着用户提供多少条查询文件就能并行起多少个抓取任务。6.3 资源申请、限制与自动扩缩容resourceRequests/resourceLimits与autoscaling都是 src/pps/pps.proto 中管道定义的正式字段请求值ResourceSpec resource_requestsL304是 Kubernetes 调度的最低保证限制值ResourceSpec resource_limitsL305是容器的资源上限bool autoscaling 33;L474则控制无任务时是否缩容至零。示例中 GPU 通过nvidia.com/gpu资源类型声明这要求集群已安装 NVIDIA device plugin——如果你的集群没有 GPU请把resourceLimits中的 gpu 段移除训练与推理仍可在 CPU 上运行只是慢很多。七、延伸如何把这个示例改造成自己的项目这套示例的通用价值在于它演示了 ML 流水线的标准分层方式数据采集层tweets把外部 API 的拉取结果转成统一格式的文本语料注入训练所需的结构化标记模型层train加载预训练权重 → 微调 → 产出可复用模型模型本身作为 PFS 中的版本化产物被管理推理层generate消费最新模型并批量产出结果结果同样进入版本化仓库供下游消费。替换场景时只需修改对应脚本与输入输出目录约定例如把推文换成评论、新闻摘要或代码片段把twitterscraper换成任意数据源再按需调整model_name、steps、nsamples等训练与生成参数即可。得益于 Pachyderm 的版本化与自动触发每次上游数据更新都会自动产生一条新的训练 生成链路模型效果随数据迭代而持续演进。参考资料仓库内示例主文档examples/ml/gpt-2/README.md抓取脚本examples/ml/gpt-2/tweets.py管道规范examples/ml/gpt-2/tweets.pipeline.json训练脚本examples/ml/gpt-2/train.py管道规范examples/ml/gpt-2/train.pipeline.json生成脚本examples/ml/gpt-2/generate.py管道规范examples/ml/gpt-2/generate.pipeline.json镜像与部署examples/ml/gpt-2/Dockerfile、examples/ml/gpt-2/Makefile管道与输入定义src/pps/pps.proto赞分享数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载相关推荐changedetection.io Docker 部署15 分钟自建网页变更监控与通知服务changedetection.io Docker 部署15 分钟自建网页变更监控与通知服务 如果你需要盯住商品价格波动、缺货商品的补货动态或者某个官网公告后端AI 应用网页爬虫Pachyderm 分布式超参数调优实战基于 Iris 数据集与 SVM 的端到端流水线Pachyderm 分布式超参数调优实战基于 Iris 数据集与 SVM 的端到端流水线 导读 本指南基于 Pachyderm 官方示例 examples/数据工程后端云原生任务调度微服务如何快速用ip2region免费做离线IP定位从0到生产环境的完整实战如何快速用ip2region免费做离线IP定位从0到生产环境的完整实战 用户登录时你怎么判断这个IP来自哪个城市调第三方接口多一次网络往返、多一份按量计后端网络上一篇Vue-Audio-Visual未来展望音频可视化技术的发展趋势和路线图下一篇Draino进阶配置Pod保护策略与高级节点过滤规则创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表