news 2026/10/8 14:10:11

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Pachyderm 实战:利用 Cron 输入周期性从 MongoDB 摄取外部数据并写入版本化仓库
  • 数据工程
  • 后端
  • 云原生
  • 任务调度
  • 微服务

【免费下载链接】pachyderm

Data-Centric Pipelines and Data Versioning

项目地址:https://gitcode.com/gh_mirrors/pa/pachyderm
点击查看免费下载

本指南以 Pachyderm 仓库中的 examples/db/README.md 示例为主线,完整演示如何在 Pachyderm 集群之外运行一个 MongoDB 实例,并借助 Pachyderm 的cron输入类型定时执行查询、将结果写入版本化的输出仓库。阅读完本文,你将掌握 MongoDB Atlas 的初始化与数据导入、通过 Kubernetes Secret(kubectl/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,按以下步骤操作:

  1. 部署一个名为Cluster0的新集群,务必记住管理员用户名与密码,后续连接与管道凭据都会用到。部署完成后可在 Atlas 仪表盘中看到该集群(如上图)。

  2. 点击集群的 "connect" 按钮,将IP 白名单配置为所有 IP(0.0.0.0/0),或者至少放行 Pachyderm 所在 Kubernetes 集群的 master 节点 IP。否则管道中的mongo容器将无法建立连接。

  1. 选择 "Connect with the MongoDB shell",记录下连接URI、数据库名(AtlasCluster0默认是test)、用户名以及认证数据库(authentication DB),这些将用于查询 MongoDB。

  2. 在本地安装 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 MB,25359 条文档)后,用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 传递这五个键值:

  • uri
  • username
  • password
  • db
  • collection

3.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.2(Kubernetes 路径)用 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.3(Pachyderm 路径)用 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 中的CreateSecretRPC(src/pps/pps.proto),客户端实现为APIClient.CreateSecret(src/client/pps.go),请求体直接携带完整 Secret 文件内容,由服务端写入 Kubernetes。Secret消息体定义在 src/pps/pps.proto,包含name与data(base64 键值映射),这正是mongodb-credentials-template.jq所构造的字段结构。

第四步:编写并创建定时查询管道

完整的管道规范见仓库中的 query.pipeline.json,它完成三件事:

  1. 挂载 Secret:transform.secrets声明名为mongosecret的 Secret,挂载到/tmp/mongosecret;
  2. 定义定时输入:input.cron每 10 秒触发一次(spec: "@every 10s");
  3. 执行查询:基于官方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),它支持四种字段:

  • name:Kubernetes 中的 Secret 名称;
  • key:Secret 中某个键,仅当同时设置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:输入名称(示例为tick);
  • repo/commit:cron 输入对应的仓库与提交;
  • spec:cron 表达式(示例为@every 10s);
  • overwrite:为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循环读取query@master:output.json,即可实时看到每次查询随机抽到的不同餐馆文档:

$ watch pachctl get file query@master:output.json

也可以按提交 ID 追溯任意历史时刻的结果——这正是 Pachyderm 数据版本化的直接体现,每个提交都能精确还原当时查询到的内容:

$ pachctl get file query@master: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 query@64ac2bd721d04212a3a0b90833f751e5: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 query@74a6cf68de2047fe94ac7982065df03d: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 测试集群中执行;文中示例截图与输出来自示例文档演示环境,实际输出内容与时间戳会随你的集群与数据有所不同。

  • 数据工程
  • 后端
  • 云原生
  • 任务调度
  • 微服务

【免费下载链接】pachyderm

Data-Centric Pipelines and Data Versioning

项目地址:https://gitcode.com/gh_mirrors/pa/pachyderm
点击查看免费下载

相关推荐

上一篇:MagiskOnWSA终极指南:在Windows上构建完整Android开发环境
下一篇:5大核心优势:Python评分卡开发终极指南 - 从零构建金融风控模型的高效方案

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/8 14:02:09

合伙人管理软件系统研发:流量资本化与权益分配机制拆解

做管理软件研发这些年&#xff0c;我发现一个很有意思的现象&#xff1a;大部分合伙人制度死在“软件太简单”上。你以为签个协议、开个账户、按比例分红就完事了&#xff0c;但真正做起来才发现&#xff0c;最难的往往不是分钱&#xff0c;而是怎么定义“流量贡献”&#xff0…

作者头像 李华
网站建设 2026/10/8 14:01:06

2026 想拿 AI 工程师 Offer?别只会调 API——这 8 层能力才是分水岭

摘要很多人以为懂 AI 就是会调用 LLM API、会写几句 prompt。2026 年&#xff0c;这个门槛已经过时了。真正值钱的 AI 工程师&#xff0c;是能把模型放进生产系统里跑起来、还能在它崩的时候救回来的人。一张流传很广的《2026 AI Engineer 通关路线图》&#xff0c;把能力拆成了…

作者头像 李华
网站建设 2026/10/8 14:00:26

ARM64离线部署Tendis单机版:镜像与compose全链路解析

简介&#xff1a;面向ARM64架构CPU环境下离线部署Tendis 2.7.0单机版的运维与开发人员&#xff0c;这份工具包以docker-compose一键脚本方式完成部署、启动、停止、卸载与健康检测&#xff0c;特别适合内网隔离或无法访问外网镜像的生产环境。包内共19个文件&#xff0c;类型涵…

作者头像 李华
网站建设 2026/10/8 13:53:43

基于人脸表情识别的课堂行为检测实战:从数据到专注度评分

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华