资讯动态

Pachyderm Word Count 实战指南:基于 MapReduce 管线的分布式词频统计与数据版本化

发布时间:2026/10/9 3:10:46 来源:尧图企业网站定制
数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载本文档以仓库 examples/word_count/README.md 为主体结合scraper、map、reduce三条管线定义pipelines、map 阶段源码 src/map.go 及配套 Makefile 展开。以下内容仅介绍查看、安装、运行与配置方式仓库为只读资源无需也不应修改任何文件。Pachyderm 的 word count 示例是一套经典的 MapReduce 入门实战用一条scraper管线抓取网页、一条map管线并行分词并统计每个词在单页的出现次数、再由reduce管线聚合出全站词频。读完本文你将掌握如何用pachctl创建输入仓库与三条管线、如何观察数据在仓库间流转、如何利用 Pachyderm 的增量处理能力追加数据并自动触发下游计算以及glob模式与并行度如何影响分发粒度。1. 示例背景与核心概念1.1 经典 MapReduce 思想在 Pachyderm 中的映射MapReduce 的核心是把输入数据切分为独立分片由map阶段并行处理再将 map 输出交给reduce阶段做聚合。在本示例中这一思想被原样映射为 Pachyderm 的三条管线Map从给定网页中提取出现过的每一个单词为每个单词生成一个以其名字命名的文本文件文件中的每一行代表该单词在某一个页面上的一次出现。Reduce把上述每个单词的逐页计数累加得到该单词在所有页面中的总出现次数。Pachyderm 的数据版本化特性让这两步天然可重放、可追溯每一次对输入仓库的提交都对应一组确定的任务所有中间结果都保存在版本化的仓库中。1.2 本示例涉及的三个关键概念在动手前先理解本示例依赖的三个 Pachyderm 概念本仓库 README 将其列为前置知识文件追加策略File Appending把文件放入 Pachyderm 仓库时如果同名文件已存在Pachyderm默认将新数据追加到已有文件末尾除非显式添加override标志。这正是本示例第二次加入Github页面后license文件从 2 行变为 3 行的原因见第 4 节。并行度Parallelism与 Glob 模式Glob Patternglob决定如何把输入仓库的数据切分成一个个 datum进而决定并行分发的粒度。三条管线的输入分别使用了/*、/*/*、/三种粒度见第 2 节逐条解读。pipeline 与 repo 的关系每条管线的输入是某个仓库的一个文件集合输出则自动落在一个与管线同名的输出仓库中层层串联形成数据流。2. 准备工作Getting Ready2.1 前置条件一个运行中的 Pachyderm 工作环境可参考本地安装方式部署。安装pachctl命令行工具并已创建好 context即已登录到集群。2.2 验证环境连通克隆本仓库后先确认 Pachyderm 集群可用$ pachctl version COMPONENT VERSION pachctl 1.12.0 pachd 1.12.0理想情况下pachctl与pachd的版本应一致最低要求是保持 major 与 minor 版本相同例如同为 1.12.x。README 同时提示每个 Pachyderm 小版本都可能带来较大的架构变化因此本仓库按分支维护不同版本的示例——master 分支对应 2.1.x另有 2.0.x、1.13.x 分支。2.3 直接运行或自建镜像本示例可直接按原文运行你也可以自行构建、打标签并推送镜像到自己的 Docker Hub。若要自建镜像需要修改 Makefile 顶部的CONTAINER_TAG默认为pachyderm/example-wordcount:2.0.3同步更新管线定义中引用的镜像执行make docker-image构建并推送镜像。Makefile 中的镜像构建目标使用docker buildx同时构建 amd64 与 arm64 平台镜像再用docker manifest create/push合并为多架构 manifest便于在异构节点上运行。3. 管线结构三个串联的处理阶段本示例由三条首尾相接的管线组成图片 展示了完整的数据流与仓库流转3.1 输入仓库urlsurls是入口仓库用于存放包含 URL 的文件。每个文件名代表要抓取的站点文件内容是该站点的具体页面地址。本示例的 data/Wikipedia 内容为https://en.wikipedia.org/wiki/Main_Page https://en.wikipedia.org/wiki/Pachyderm即针对 Wikipedia 站点抓取两个页面。3.2 三条管线的分工管线定义文件输入输出仓库作用scraperscraper.pipeline.jsonurlsglob: /*scraper抓取 URL 对应的 .html 内容mapmap.pipeline.jsonscraperglob: /*/*map并行分词按词输出逐页计数reducereduce.pipeline.jsonmapglob: /reduce聚合每个词的总计数三条管线含reduce都可以分布式运行以最大化性能——只要把输入仓库的glob切分成多个 datumPachyderm 就会把不同 datum 分发到不同 worker 并行处理。3.3 逐条解读管线定义scraper管线scraper.pipeline.json使用基础镜像执行一段 bash 脚本安装wget后遍历/pfs/urls/*对每个 URL 文件调用wget抓取页面输出目录按源文件名$(basename $f)分目录存放glob: /*表示每个 URL 文件作为一个 datum两个 URL 文件可并行抓取。定义中还设置了acceptReturnCode: [4, 5, 6, 7, 8]即把 wget 常见的网络错误返回码视为可接受的不导致任务失败增强抓取健壮性。map管线map.pipeline.json运行pachyderm/example-wordcount:2.0.3镜像中的app可执行程序参数为/pfs/scraper/与/pfs/outglob: /*/*表示每个 HTML 文件位于scraper/站点/页面.html是一个 datum从而对每页并行执行分词。其 Go 实现见 src/map.go后面详述。reduce管线reduce.pipeline.json输入glob: /即把map输出仓库整体作为一个 datum在单任务内用 shell 把/pfs/map/*/*下每个单词文件的所有行合并再用awk { sum$1 } END { print sum }累加得到总计数写入/pfs/out/词。4. 示例第一部分跑通端到端词频统计Step 1创建入口仓库并建立三条管线进入examples/word_count目录最快的方式是直接执行$ make wordcount该目标见 Makefile等价于依次执行以下命令$ pachctl create repo urls $ cd data pachctl put file urlsmaster -f Wikipedia其中 data/Wikipedia 含 2 个 URL指向 2 个 Wikipedia 页面随后创建三条管线$ pachctl create pipeline -f pipelines/scraper.pipeline.json $ pachctl create pipeline -f pipelines/map.pipeline.json $ pachctl create pipeline -f pipelines/reduce.pipeline.json说明Makefile 中的命令带有--project $(PROJECT)参数默认PROJECT ? default即按 Pachyderm 的 project 作用域创建手工执行时若使用默认 project 可省略。此外 Makefile 还提供了clean目标可一次性删除三条管线与四个仓库用于重跑。由于第一条管线的入口仓库urls中已存在待处理数据创建管线会立即触发 3 个任务。用以下命令观察任务执行$ pachctl list job输出形如ID PIPELINE STARTED DURATION RESTART PROGRESS DL UL STATE ba77c94678ae401db1a2b58528b74d78 reduce 42 seconds ago - 0 0 0 / 1488 0B 0B running 7e2ed8e0a8dd49a18f1e45d62942c7ee map 43 seconds ago Less than a second 0 2 0 / 2 5.218MiB 4.495KiB success bc3031c8076c4f91aa1d8e6ba450b096 map_build 50 seconds ago 2 seconds 0 1 0 / 1 1.455KiB 2.557MiB success badd3d81d3ce46358d91bedbb34dd0ed scraper 59 seconds ago 16 seconds 0 1 0 / 1 81B 103.4KiB success注意其中出现的map_build它是 Pachyderm 为 map 管线自动执行的镜像构建任务把 Dockerfile 打包为app可执行程序完成后才运行真正的 map 任务。随后检查仓库是否全部创建$ pachctl list repoNAME CREATED SIZE (MASTER) DESCRIPTION reduce 13 minutes ago 4.62KiB Output repo for pipeline reduce. map 13 minutes ago 4.495KiB Output repo for pipeline map. scraper 13 minutes ago 103.4KiB Output repo for pipeline scraper. urls 13 minutes ago 81B四个仓库的流转关系为urls→scraper→map→reduce每个仓库都有对应的DESCRIPTION说明其来源管线。Step 2逐层检查仓库内容scraper 仓库——查看抓取结果$ pachctl list file scrapermasterNAME TYPE SIZE /Wikipedia dir 103.4KiB$ pachctl list file scrapermaster:/WikipediaNAME TYPE SIZE /Wikipedia/Main_Page.html file 77.53KiB /Wikipedia/Pachyderm.html file 25.88KiB两个 URL 成功抓取为两个 .html 页面。map 仓库——查看分词结果。map 管线为每个单词生成一个以其名字命名的文件实现见 src/map.go$ pachctl list file mapmasterNAME TYPE SIZE ... /liberated file 2B /librarians file 2B /library file 2B /license file 4B /licenses file 4B ...每个文件内容是该词在每个页面上的出现次数一行一页。查看license一词$ pachctl get file mapmaster:/license 5 5即单词license在两个页面中各出现 5 次。默认情况下Pachyderm 会为集群中的每个节点各启动一个 worker可通过管线定义或相关配置调整 worker 数量见 README 引用的 Distributed Computing 文档。reduce 仓库——查看聚合后的总词频$ pachctl list file reducemasterNAME TYPE SIZE ... /liberated file 2B /librarians file 2B /library file 2B /license file 3B /licenses file 2B ...验证license的总计数是否为 55$ pachctl get file reducemaster:/license 10reduce输出文件3B比map阶段的对应文件4B小是因为它只保留了一行累加结果而 map 阶段保留了逐页的每一行。5. 理解 map 阶段的分词实现源码解读src/map.go 是 map 管线的核心程序约 90 行 Go 代码由 Dockerfile 基于golang:1.17.6构建为/go/bin/app。其处理流程为参数解析程序接收两个命令行参数——输入目录/pfs/scraper/与输出目录/pfs/out个数不符则log.Fatalf退出。分词与清洗编译正则[^A-Za-z]把扫描到的每个词中非字母字符替换为空格再转小写并切分得到干净的英文小写单词列表。遍历统计用filepath.Walk递归遍历输入目录对每个文件用bufio.Scanner按单词bufio.ScanWords扫描累加每个单词的出现次数日志会输出每个文件扫描到的总词数found %d words in %s。按词输出遍历内存中的wordMap把每个单词写入输出目录下以PACH_DATUM_ID环境变量命名的子目录中的同名文件内容为该词在本 datum 内的计数一行一个数字。从实现细节可以看到由于glob: /*/*使每个 HTML 文件成为一个 datumPachyderm 会为每个 datum 注入独立的PACH_DATUM_IDmap 程序据此把结果写入互不冲突的子目录从而支持并行执行一个 datum 内只统计一个页面因此mapmaster:/license的一行即一个页面的计数reduce阶段把map仓库所有 datum 的结果按文件名合并累加得到全量总计数。6. 示例第二部分增量添加数据观察管线自动扩算端到端流程跑通后再向urls仓库追加一个新站点。得益于 Pachyderm 的文件追加语义新增文件不会重跑旧数据$ cd data pachctl put file urlsmaster -f Githubdata/Github 内容为https://github.com/pachyderm/pachydermscraper管线会自动开始抓取新站点不会重新抓取 Wikipedia随后自动触发map、reduce处理新数据并把所有站点的词频合并更新。6.1 scraper 仓库新增了一个站点目录$ pachctl list file scrapermasterNAME TYPE SIZE /Github dir 195.1KiB /Wikipedia dir 103.4KiB$ pachctl list file scrapermaster:/GithubNAME TYPE SIZE /Github/pachyderm.html file 195.1KiB6.2 map 仓库出现新词、旧词文件增长$ pachctl list file mapmaster对比新增前后新出现了单词libraries的文件license文件大小从 4B 变为 7B多了一行$ pachctl get file mapmaster:/license 23 5 5说明单词license在新抓取的 Github 页面上出现了 23 次。注意这里正是第 1.2 节提到的文件追加策略map 输出仓库中已存在的/license文件没有被打上override标志因此新 datum 的计数被追加到文件末尾而不是覆盖原有两行。6.3 reduce 仓库总计数自动更新$ pachctl list file reducemasterreduce输出目录中新增了libraries文件license文件仍为 3B一行总计数。验证 2355 是否等于 33$ pachctl get file reducemaster:/license 33reduce 管线把 map 阶段追加的多行计数重新聚合成一行最终词频统计始终保持全量一致。7. 调优提示并行度与 Glob 模式README 末尾给出了两条重要的调优指引worker 数量默认情况下每条管线会为集群中的每个节点各启动一个 worker你可以在管线定义中设置不同的 worker 数量参考 Distributed Computing 文档。glob 模式本示例三条管线已配置好分发粒度——scraper用glob: /*按 URL 文件分发、map用glob: /*/*按页面分发、reduce用glob: /整体聚合。这正体现了 glob 模式对性能的直接影响粒度越细并行度越高但聚合代价也越大实践中应根据数据规模和计算成本权衡选择。8. 小结与下一步通过本示例你可以完整观察到 Pachyderm 数据管线的核心工作方式仓库即数据流urls→scraper→map→reduce四个仓库层层串联每条管线的输出自动成为下一条管线的输入声明式管线三条 pipeline JSON 只声明输入、命令与镜像无需关心任务调度细节增量处理追加新输入data/Github时旧数据不会被重算新数据自动触发下游更新可验证结果从pachctl list job、pachctl list file、pachctl get file可以全程追踪每个阶段的输入输出与最终词频。如果想继续深入可以进一步阅读 examples/word_count/src/map.go 了解分词细节或者对照本仓库其他示例如 examples/joins、examples/group学习 join、group 等更高级的管线操作。赞分享数据工程后端云原生任务调度微服务【免费下载链接】pachydermData-Centric Pipelines and Data Versioning项目地址https://gitcode.com/gh_mirrors/pa/pachyderm点击查看免费下载相关推荐Apache Beam 入门实战Java Kata 之 Word Count 词频统计管线解析Apache Beam 入门实战Java Kata 之 Word Count 词频统计管线解析 导读 Word Count词频统计是 Apache Bea大数据批处理流处理数据工程PyFlink Table API Word Count 实战从批式 CSV 处理到 Datagen 流式词频统计PyFlink Table API Word Count 实战从批式 CSV 处理到 Datagen 流式词频统计 本文以 flink python/docs后端大数据流处理批处理Pachyderm终极实战指南5步构建高效MapReduce单词计数应用Pachyderm终极实战指南5步构建高效MapReduce单词计数应用 Pachyderm是一个强大的分布式数据仓库和数据处理平台专为大规模数据分析和机器数据工程后端云原生任务调度微服务上一篇Algorithm Visualizer 项目常见问题解决方案下一篇Parse Server 常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑