ARTICLE DETAIL

资讯详情

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

Pachyderm 实战:利用 Cron 输入周期性从 MongoDB 摄取外部数据并写入版本化仓库

Pachyderm 实战:利用 Cron 输入周期性从 MongoDB 摄取外部数据并写入版本化仓库 数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载本指南以 Pachyderm 仓库中的 examples/db/README.md 示例为主线完整演示如何在 Pachyderm 集群之外运行一个 MongoDB 实例并借助 Pachyderm 的cron输入类型定时执行查询、将结果写入版本化的输出仓库。阅读完本文你将掌握 MongoDB Atlas 的初始化与数据导入、通过 Kubernetes Secretkubectl/pachctl两种路径向管道安全传递数据库凭据、编写含cron输入的管道规范以及用pachctl观察周期任务与按提交追溯历史结果的完整实战技能。示例背景Pachyderm 与外部数据库的周期数据摄取Pachyderm 的定位是数据为中心的流水线与数据版本化Data-Centric Pipelines and Data Versioning。多数示例聚焦于对仓库内已有数据做转换但真实生产场景中经常需要从Pachyderm 之外的数据库、消息队列或 SaaS 服务周期性拉取数据。本示例正好补上这一环一条名为query的管道每隔 10 秒对集群外部的 MongoDB 执行一次$sample聚合查询随机抽取一条餐馆记录写入query输出仓库。由于每次查询都会产生一个新的提交commit输出仓库天然形成随时间演进的版本化数据流可被下游管道周期性消费。示例实现要点如下通过管道transform.secrets将 MongoDB 的连接 URI、账号、密码、库名、集合名以 Kubernetes Secret 方式挂载到容器内通过input.cron定义定时触发策略every 10s无需上游数据提交也能驱动任务使用官方mongo镜像在管道内直接执行查询结果写入/pfs/out/output.json。运行本示例前你需要具备一个正在运行的 Pachyderm 集群可使用官方 Local Installation 方式在本地几分钟内启动安装并已连接到该集群的pachctl命令行工具。第一步在 MongoDB Atlas 上准备 MongoDB 集群最省事的方式是使用免费的托管 MongoDB 服务如 MongoDB Atlas 的免费层当然也可以使用任何你能访问的 MongoDB 实例。若使用 MongoDB Atlas按以下步骤操作部署一个名为Cluster0的新集群务必记住管理员用户名与密码后续连接与管道凭据都会用到。部署完成后可在 Atlas 仪表盘中看到该集群如上图。点击集群的 connect 按钮将IP 白名单配置为所有 IP0.0.0.0/0或者至少放行 Pachyderm 所在 Kubernetes 集群的 master 节点 IP。否则管道中的mongo容器将无法建立连接。选择 Connect with the MongoDB shell记录下连接URI、数据库名AtlasCluster0默认是test、用户名以及认证数据库authentication DB这些将用于查询 MongoDB。在本地安装 MongoDB 命令行工具如mongoimport、mongoshell后续导入数据集与调试查询都会用到。第二步导入示例数据到 MongoDB本示例使用 MongoDB 官方示例数据集primer-dataset.json内容是纽约市的餐馆记录也是 MongoDB 官方文档中反复使用的经典数据集。每条记录的字段结构如下{ address: { building: 1007, coord: [ -73.856077, 40.848447 ], street: Morris Park Ave, zipcode: 10462 }, borough: Bronx, cuisine: Bakery, grades: [ { date: { $date: 1393804800000 }, grade: A, score: 2 }, { date: { $date: 1378857600000 }, grade: A, score: 6 }, { date: { $date: 1358985600000 }, grade: A, score: 10 }, { date: { $date: 1322006400000 }, grade: A, score: 9 }, { date: { $date: 1299715200000 }, grade: B, score: 14 } ], name: Morris Park Bake Shop, restaurant_id: 30075445 }下载该数据集文件名为primer-dataset.json约 11.3 MB25359 条文档后用mongoimport将其导入test库Atlas 默认库名的restaurants集合。命令中需要按你的集群实际情况填写主机列表、用户名、密码与认证数据库$ mongoimport --host Cluster0-shard-0/cluster0-shard-00-00-cwehf.mongodb.net:27017,cluster0-shard-00-01-cwehf.mongodb.net:27017,cluster0-shard-00-02-cwehf.mongodb.net:27017 --ssl -u admin -p my password --authenticationDatabase admin --db test --collection restaurants --drop --file primer-dataset.json 2017-08-28T13:40:38.983-0400 connected to: Cluster0-shard-0/cluster0-shard-00-00-cwehf.mongodb.net:27017,cluster0-shard-00-01-cwehf.mongodb.net:27017,cluster0-shard-00-02-cwehf.mongodb.net:27017 2017-08-28T13:40:39.048-0400 dropping: test.restaurants ... 2017-08-28T13:42:08.449-0400 [########################] test.restaurants 11.3MB/11.3MB (100.0%) 2017-08-28T13:42:08.449-0400 imported 25359 documents其中--drop会先清空同名集合再导入便于重复执行导入成功后输出imported 25359 documents即表示restaurants集合已就绪。第三步将 MongoDB 凭据封装为 Kubernetes Secret管道需要知道 MongoDB 的 URI、用户名、密码、库名与集合名。示例通过 Kubernetes Secret 传递这五个键值uriusernamepassworddbcollection3.1 将凭据写入本地文件先把各值写入本地文件值必须用单引号包裹防止 shell 解释特殊字符并用chmod 600收紧权限$ echo -n uri uri ; chmod 600 uri $ echo -n username username ; chmod 600 username $ echo -n password password ; chmod 600 password $ echo -n db db ; chmod 600 db $ echo -n collection collection ; chmod 600 collection创建后逐一确认内容无误$ cat uri $ cat username $ cat password $ cat db $ cat collection创建 Secret 有两种路径有 Kubernetes 直接访问权限时用kubectl下文 Kubernetes 路径没有或不想用kubectl时走pachctl下文 Pachyderm 路径。3.2Kubernetes 路径用 kubectl 创建 Secret$ kubectl create secret generic mongosecret --from-file./uri \ --from-file./username \ --from-file./password \ --from-file./db \ --from-file./collection验证方式是把 Secret 导出为 JSON 并用jq的base64d解码核对$ kubectl get secret mongosecret -o json | jq .data | map_values(base64d) { uri: uri, username: username password: password db: db collection: collection }3.3Pachyderm 路径用 pachctl 创建 Secret先用仓库中提供的 mongodb-credentials-template.jq 模板把五个文件的值经 base64 编码后组装成 Kubernetes Secret 的 JSON 定义$ jq -n --arg uri $(cat uri) --arg username $(cat username) \ --arg password $(cat password) --arg db $(cat db) --arg collection $(cat collection) \ -f mongodb-credentials-template.jq mongodb-credentials-secret.json $ chmod 600 mongodb-credentials-secret.json解码核对生成文件的内容$ jq .data | map_values(base64d) mongodb-credentials-secret.json { uri: uri, username: username password: password db: db collection: collection }最后通过 pachctl 创建 Secret$ pachctl create secret -f mongodb-credentials-secret.json从源码看pachctl create secret底层调用的是 PPS API 中的CreateSecretRPCsrc/pps/pps.proto客户端实现为APIClient.CreateSecretsrc/client/pps.go请求体直接携带完整 Secret 文件内容由服务端写入 Kubernetes。Secret消息体定义在 src/pps/pps.proto包含name与database64 键值映射这正是mongodb-credentials-template.jq所构造的字段结构。第四步编写并创建定时查询管道完整的管道规范见仓库中的 query.pipeline.json它完成三件事挂载 Secrettransform.secrets声明名为mongosecret的 Secret挂载到/tmp/mongosecret定义定时输入input.cron每 10 秒触发一次spec: every 10s执行查询基于官方mongo镜像随机抽取restaurants集合中的一条文档并写入/pfs/out。{ pipeline: { name: query }, transform: { image: mongo, cmd: [ /bin/bash ], stdin: [ export uri$(cat /tmp/mongosecret/uri), export db$(cat /tmp/mongosecret/db), export collection$(cat /tmp/mongosecret/collection), export username$(cat /tmp/mongosecret/username), export password$(cat /tmp/mongosecret/password), mongo \$uri\ --authenticationDatabase admin --ssl --username $username --password $password --quiet --eval db.restaurants.aggregate({ $sample: { size: 1 } }); | tail -n1 | egrep -v \^|^bye\ /pfs/out/output.json ], secrets: [ { name: mongosecret, mount_path: /tmp/mongosecret } ] }, input: { cron: { name: tick, spec: every 10s } } }4.1 Secret 挂载的底层机制transform.secrets的类型对应 PPS proto 中的SecretMount消息src/pps/pps.proto它支持四种字段nameKubernetes 中的 Secret 名称keySecret 中某个键仅当同时设置env_var时才有意义用于以环境变量方式注入mount_path将 Secret 挂载为文件系统的目标路径env_var可选若设置则把对应键的值写入指定的环境变量。本示例采用mount_path: /tmp/mongosecret的文件挂载方式因此容器内/tmp/mongosecret/uri、/tmp/mongosecret/db等路径就是 Secret 中各键的内容stdin里的cat命令正是读取这些文件。worker 侧由 src/server/pps/server/worker_rc.go 负责把 Secret 转成 Kubernetes volume 与 volumeMount 注入到用户容器。4.2 Cron 输入的语义CronInput消息定义在 src/pps/pps.proto核心字段包括name输入名称示例为tickrepo/commitcron 输入对应的仓库与提交speccron 表达式示例为every 10soverwrite为true时每次 tick 覆盖同一个 datum为false时每个 tick 创建新 datum默认行为即本例start可选的起始时间戳决定何时开始调度。cron 表达式的解析由 src/internal/cronutil/cronutil.go 的ParseCronExpression完成它包装了 robfig/cron 库的cron.ParseStandard因此既支持标准的 5 段 cron 语法* * * * *也支持every 10s这类描述式写法。仓库配套的单测 src/internal/cronutil/cronutil_test.go 覆盖了every 1m等表达式的解析验证。理解cron输入的关键是它不依赖上游数据而是由调度器在每个 tick 主动生成输入。每个 tick 对应一个提交与一个 datum管道随之触发一次 job——这正是本示例能周期性从外部数据库拉数据的根本原因。创建管道$ pachctl create pipeline -f query.pipeline.json第五步观察任务触发与版本化结果创建后用pachctl list pipeline确认管道状态INPUT列会显示 cron 输入及其调度表达式$ pachctl list pipeline NAME VERSION INPUT CREATED STATE / LAST JOB DESCRIPTION query 1 tick:every 10s 6 seconds ago running / starting管道启动后每 10 秒应能看到一个新 job 被触发。连续执行pachctl list job可看到running状态的新任务与一系列success的历史任务交错出现每个 job 的OUTPUT COMMIT对应一个独立提交 ID$ pachctl list job ID OUTPUT COMMIT STARTED DURATION RESTART PROGRESS DL UL STATE 5938a0d0-9512-455f-a390-14adc3669e5f query/0f8a2ba1150a463299ee71961427bdcb 3 seconds ago 3 seconds 0 1 0 / 1 26B 617B success 952427a6-c92d-4c98-a781-87616988d528 query/33776e4df3b24ab68d70b5185eb37661 13 seconds ago 1 second 0 1 0 / 1 26B 613B success 1bc5f608-85fd-44eb-833e-562d15629706 query/6dd2a4da566f4d30ad9c66fc60244bab 23 seconds ago 1 second 0 1 0 / 1 26B 721B success efa677a4-7f83-424b-879d-70a0c5690bb2 query/f56b1f314030455c8bdf8a10b68ebd16 33 seconds ago 1 second 0 1 0 / 1 26B 529B success 842e4e6c-4920-42c0-9c81-e5299b67e4a0 query/2a11bfc3e6d74af0a8d254d3ecf6f6af 43 seconds ago 1 second 0 1 0 / 1 26B 535B success其中每个 job 的DL下载 26B对应 cron 输入产生的 tick 数据UL数百字节不等则是查询结果output.json的写入量。用watch循环读取querymaster:output.json即可实时看到每次查询随机抽到的不同餐馆文档$ watch pachctl get file querymaster:output.json也可以按提交 ID 追溯任意历史时刻的结果——这正是 Pachyderm 数据版本化的直接体现每个提交都能精确还原当时查询到的内容$ pachctl get file querymaster:output.json { _id : ObjectId(59a455af69a077c0dc028410), address : { building : 119, coord : [ -73.9784962, 40.6788476 ], street : 5 Avenue, zipcode : 11217 }, borough : Brooklyn, cuisine : Mexican, grades : [ { date : ISODate(2014-07-29T00:00:00Z), grade : B, score : 27 }, ... ], name : El Pollito Mexicano, restaurant_id : 41051406 } $ pachctl get file query64ac2bd721d04212a3a0b90833f751e5:output.json { _id : ObjectId(59a455f069a077c0dc02e16e), address : { building : 1650, coord : [ -73.928079, 40.856481 ], street : Saint Nicholas Ave, zipcode : 10040 }, borough : Manhattan, cuisine : Spanish, grades : [ { date : ISODate(2015-01-20T00:00:00Z), grade : Not Yet Graded, score : 2 } ], name : Angebienvendia, restaurant_id : 50018661 } $ pachctl get file query74a6cf68de2047fe94ac7982065df03d:output.json { _id : ObjectId(59a455b669a077c0dc02904d), address : { building : 14, coord : [ -73.990382, 40.741571 ], street : West 23 Street, zipcode : 10010 }, borough : Manhattan, cuisine : Café/Coffee/Tea, grades : [ { date : ISODate(2014-05-02T00:00:00Z), grade : A, score : 11 }, ... ], name : Starbucks Coffee (Store #13539), restaurant_id : 41290548 }扩展思考从周期查询到下游流水线本示例的输出仓库query本身可以继续作为下游管道的输入。例如你可以追加一条管道以query仓库为pfs输入对每个新提交做清洗、聚合或入库从而把外部数据 → Pachyderm → 结果仓库的周期链路扩展为多级数据流。Pachyderm 的提交追踪provenance机制会自动记录query仓库每个提交与下游 job 的依赖关系保证整条链路的可复现性。若想调整查询频率只需修改query.pipeline.json中的input.cron.spec例如every 1m或标准 cron 表达式*/5 * * * *再执行pachctl update pipeline即可生效CronInput的overwrite字段则决定了每个 tick 是追加新 datum 还是覆盖旧 datum适用于只关心最新快照的场景。小结本示例展示了 Pachyderm 接入外部数据库的完整范式Secret 传递凭据、cron 输入驱动周期任务、输出仓库累积版本化结果。配套的 query.pipeline.json 与 mongodb-credentials-template.jq 均可直接复制使用只需替换为你的 MongoDB 连接信息。需要注意的是Pachyderm 各 minor 版本间示例可能随架构演进调整仓库将不同版本的示例维护在对应分支中本仓库 master 分支对应的示例基于 Pachyderm 2.1.x 系列。说明运行本示例涉及创建/更新管道、Secret 等集群操作请仅在你自己的 Pachyderm 测试集群中执行文中示例截图与输出来自示例文档演示环境实际输出内容与时间戳会随你的集群与数据有所不同。赞分享数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载相关推荐DataHub Power BI 元数据摄取实战从 Entra 应用注册到周期性摄取管线DataHub Power BI 元数据摄取实战从 Entra 应用注册到周期性摄取管线 本文围绕 DataHub 官方的 Power BI 快速摄取指南数据目录数据治理数据血缘后端前端数据工程数据集成Pachyderm数据访问模式读取优化与写入优化策略Pachyderm数据访问模式读取优化与写入优化策略 在当今数据驱动的时代 Pachyderm 作为一款强大的分布式数据仓库和数据处理平台其数据访问模式的数据工程后端云原生任务调度微服务Claude SEO 实战FLOW Optimize 阶段的 Follow-Up Qualifying Prompt——证据筛选、优先级排序与发布前验证清单Claude SEO 实战FLOW Optimize 阶段的 Follow Up Qualifying Prompt——证据筛选、优先级排序与发布前验证清单上一篇MagiskOnWSA终极指南在Windows上构建完整Android开发环境下一篇5大核心优势Python评分卡开发终极指南 - 从零构建金融风控模型的高效方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表