5.2 数据流工具Pig与生态补充


文档摘要

5.2 数据流工具 Pig 与生态补充 本节摘要:Pig 用 Pig Latin 脚本描述数据流——逐行声明"加载、过滤、连接、分组、存储",框架自动优化并翻译成 MapReduce。本节用一个完整脚本走一遍语言核心,对比 Pig 与 Hive 的适用分界,并补齐 ZooKeeper、Mahout、Ambari 等周边组件在生态中的角色。 先写脚本,再谈概念 任务还是那个:从访问日志统计活跃用户。Pig Latin 的写法是"一步步说清数据怎么流": 读这段脚本的感觉应当是:每行就是数据旅程的一个加工站,从上往下读就是数据从原料到成品的流水线。没有子查询、没有嵌套 SQL,所有中间结果都有名字。这就是 Pig 的哲学:数据流描述语言,程序员式的直觉。 四个机制关键词 惰性执行。

5.2 数据流工具 Pig 与生态补充

本节摘要:Pig 用 Pig Latin 脚本描述数据流——逐行声明"加载、过滤、连接、分组、存储",框架自动优化并翻译成 MapReduce。本节用一个完整脚本走一遍语言核心,对比 Pig 与 Hive 的适用分界,并补齐 ZooKeeper、Mahout、Ambari 等周边组件在生态中的角色。

先写脚本,再谈概念

任务还是那个:从访问日志统计活跃用户。Pig Latin 的写法是"一步步说清数据怎么流":

-- 加载:指路径与模式 分隔符制表符 logs = LOAD '/data/access/dt=2026-08-18/' USING PigStorage('\t') AS (ts:chararray, user_id:chararray, url:chararray, bytes:int); -- 过滤:脏行与爬虫一并去掉 clean = FILTER logs BY url is not null AND NOT (url MATCHES '.*spider.*'); -- 提取需要的字段(投影) proj = FOREACH clean GENERATE user_id, bytes; -- 按用户分组 grp = GROUP proj BY user_id; -- 每组聚合:总流量与次数 agg = FOREACH grp GENERATE group AS user_id, SUM(proj.bytes) AS total_bytes, COUNT(proj) AS pv; -- 只留重度用户 heavy = FILTER agg BY pv > 100; -- 存储:一份文本 一份供后续join的有序版本 STORE heavy INTO '/out/heavy_users' USING PigStorage(','); sorted = ORDER heavy BY total_bytes DESC; STORE sorted INTO '/out/heavy_users_sorted';

读这段脚本的感觉应当是:每行就是数据旅程的一个加工站,从上往下读就是数据从原料到成品的流水线。没有子查询、没有嵌套 SQL,所有中间结果都有名字。这就是 Pig 的哲学:数据流描述语言,程序员式的直觉。

四个机制关键词

惰性执行。LOAD/FILTER/GROUP 都不触发计算,只有 STORE 触发——Pig 把整条流水线编译成一个(或少数几个)MR 作业,做全局优化:FILTER 尽可能推到最前(早过滤)、投影尽早裁列、join 策略自动选择(小表广播还是 shuffle)。换句话说,你写的是手工流水线,跑的是优化器重排过的版本——这与 SQL 声明式殊途同归。

嵌套数据类型。bag(元组的集合)、tuple(有序字段列表)、map(键值对)可以嵌套。GROUP 产生的就是"每组一个 bag",FOREACH 里还能对内层 bag 做嵌套操作(排序、过滤、取 top):

top3 = FOREACH grp { sorted_inner = ORDER proj BY bytes DESC; top = LIMIT sorted_inner 3; GENERATE group AS user_id, top; }

SQL 里这类"组内 TopN"要么窗口函数要么自连接,Pig 的嵌套块直接表达——这是它对复杂多步变换仍有一席之地的原因。

UDF 全通。所有数据进出都经过函数钩子:加载用自定义 Loader、加工用 EVAL_FUNC 继承的 Java UDF、过滤用 FILTER UDF。大量清洗逻辑是公司私有的(脱敏、字典翻译、格式归一),Pig 的 UDF 约定极简,工程团队攒一套清洗 UDF 库可以跨脚本复用多年。

多数据流复用。一个脚本里 STORE 多个结果,共同的前置步骤只算一遍——手工 MR 要跑多个作业各扫一遍输入。

Pig 还是 Hive:分界线画在哪

两者都翻译成 MR/Tez,都吃 HDFS,生态位高度重叠,但气质不同:

维度 Hive Pig
语言 SQL 声明式 数据流过程式
模式 强模式(先建表) 松模式(脚本里声明,甚至无模式)
擅长 标准分析、报表、BI 接入 一次性清洗、多步原型、非标准变换
用户 分析师 工程师

实践中的分工惯例:稳定的有报表语义的分析进 Hive 数仓;探索性、一次性的数据搬运与清洗用 Pig 脚本。Pig 松模式对脏数据友好(无模式加载后再逐字段处理),早期互联网公司用它做大规模日志 ETL 的很多;随着 Spark 提供了同样灵活但更快(内存内)的 API,Pig 的新增使用在收缩,这层演化脉络第 7 章展开。

周边组件补位:生态的黏合剂与工具箱

按"数据旅程"的视角把其余常见组件各归其位:

ZooKeeper:分布式协调服务。第 2 章 HDFS HA 的主锁选举(ZKfc)、HBase 的 Region 寻址、Oozie 的锁、Kafka(生态外延)的控制器选举——凡是"多个进程要对同一个状态达成一致"的地方都有它。数据不流经它,但旅程中每个关键切换点都由它背书。核心抽象是 ZNode 树与观察者通知:一方写节点,多方收到变更回调,配合同一事务序列号保证顺序一致。

Mahout:机器学习库。早期把分类、聚类、推荐算法实现为 MR 作业(协方差矩阵、k-means 的 MR 化正是第 3 章分治模型的经典应用),后来转向聚焦 Spark 之上的新算法栈。今天在 Hadoop 教程里它的意义更多是"机器学习如何吃 HDFS 数据"的历史注脚与 UDF/特征工程的素材库。

Ambari:集群管理平台。装、配、起、停、监控一条龙(第 6 章运维实战的主力工具之一),对旅程的贡献是把第 2 至 4 章那一堆 XML 配置变成向导与仪表盘。

HCatalog:元数据共享层。把 Hive Metastore 的表定义开放成通用接口——Pig、MR 都能按"表名"而不是"路径"读写,Pig 产出的数据直接以 Hive 表身份被查询。生态组件靠它打通各自的模式视角,避免"Pig 写的目录 Hive 不认识"的割裂。

本节要点回顾

  • Pig Latin 是过程式数据流:每行一个加工站,STORE 才触发计算,优化器全局重排;
  • 嵌套 bag 类型让组内 TopN 这类变换直接表达,UDF 钩子让私有清洗逻辑可积累复用;
  • Pig 与 Hive 分工:标准分析进 Hive,探索清洗用 Pig;Spark 正在吸收两者领地;
  • ZooKeeper 是协调背书者:HA 主锁、Region 寻址、控制器选举都靠它,数据不流经它;
  • HCatalog 让各组件共享表视角,Ambari 把配置运维界面化,Mahout 是 ML on Hadoop 的历史注脚。

下一节进入旅程里对延迟最敏感的出口:HBase 在线查询。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U