- 数据工程
- 后端
- 云原生
- 任务调度
- 微服务
【免费下载链接】pachyderm
Data-Centric Pipelines and Data Versioning
本指南以 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,按以下步骤操作:
部署一个名为
Cluster0的新集群,务必记住管理员用户名与密码,后续连接与管道凭据都会用到。部署完成后可在 Atlas 仪表盘中看到该集群(如上图)。点击集群的 "connect" 按钮,将IP 白名单配置为所有 IP(
0.0.0.0/0),或者至少放行 Pachyderm 所在 Kubernetes 集群的 master 节点 IP。否则管道中的mongo容器将无法建立连接。
选择 "Connect with the MongoDB shell",记录下连接URI、数据库名(Atlas
Cluster0默认是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 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 传递这五个键值:
uriusernamepassworddbcollection
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,它完成三件事:
- 挂载 Secret:
transform.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),它支持四种字段:
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
相关推荐
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——证据筛选、优先级排序与发布前验证清单
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考