如何在PIG中移除多列空值?过滤操作输出异常求排查
解决Pig处理2亿行数据集空值过滤及统计输出异常问题
嘿,我来帮你搞定这个问题!2亿行的数据集规模不小,先从你遇到的输出异常入手,一步步排查和修正:
一、先排查空值过滤的潜在问题
你当前的过滤逻辑只判断了is not null,但要注意Pig里的null和空字符串''是完全不同的!如果你的数据里存在空字符串(比如字段值是空的,但不是null),is not null是不会过滤掉这些行的,这会导致后续统计结果异常。
修正后的过滤条件应该同时排除null和空字符串:
values = FILTER div BY A IS NOT NULL AND A != '' AND B IS NOT NULL AND B != '' AND C IS NOT NULL AND C != '' AND D IS NOT NULL AND D != '' AND E IS NOT NULL AND E != '' AND F IS NOT NULL AND F != '' AND G IS NOT NULL AND G != '' AND H IS NOT NULL AND H != '';
二、实现你期望的统计输出格式
你想要的H 1056 U 4355 W 999 P 1000 Y 2199这种键值对输出,分两种常见场景,对应不同的实现方式:
场景1:统计每个字段值的出现次数
如果是要统计A-H字段中所有不同取值的出现次数(比如值H出现了1056次,值U出现了4355次),可以把所有字段转成统一的键值对格式再分组统计:
-- 把所有字段拆成(field_name, value)的行,方便统一统计 melted_data = FOREACH values GENERATE 'A' AS field, A AS val UNION ALL 'B' AS field, B AS val UNION ALL 'C' AS field, C AS val UNION ALL 'D' AS field, D AS val UNION ALL 'E' AS field, E AS val UNION ALL 'F' AS field, F AS val UNION ALL 'G' AS field, G AS val UNION ALL 'H' AS field, H AS val; -- 按值分组,统计出现次数 value_counts = GROUP melted_data BY val; final_result = FOREACH value_counts GENERATE group AS value, COUNT(melted_data) AS count; -- 输出结果,Pig默认用空格分隔每行的两个值 DUMP final_result;
如果要把所有结果拼成一行,可以用shell命令后续处理:比如把输出文件导出后执行cat result.txt | tr '\n' ' '就能得到你想要的单行空格分隔格式。
场景2:统计每个字段的非空总记录数
如果你是要统计A-H每个字段各自的非空记录数(比如A字段非空有X条,B字段非空有Y条),那其实不需要先过滤所有字段非空,直接在原数据集上统计更准确:
-- 给每个字段标记是否非空(1=非空,0=空) field_flags = FOREACH dataset GENERATE (A IS NOT NULL AND A != '') ? 1 : 0 AS cnt_A, (B IS NOT NULL AND B != '') ? 1 : 0 AS cnt_B, (C IS NOT NULL AND C != '') ? 1 : 0 AS cnt_C, (D IS NOT NULL AND D != '') ? 1 : 0 AS cnt_D, (E IS NOT NULL AND E != '') ? 1 : 0 AS cnt_E, (F IS NOT NULL AND F != '') ? 1 : 0 AS cnt_F, (G IS NOT NULL AND G != '') ? 1 : 0 AS cnt_G, (H IS NOT NULL AND H != '') ? 1 : 0 AS cnt_H; -- 汇总所有标记的和,得到每个字段的非空总数 total_counts = FOREACH (GROUP field_flags ALL) GENERATE SUM(field_flags.cnt_A) AS count_A, SUM(field_flags.cnt_B) AS count_B, SUM(field_flags.cnt_C) AS count_C, SUM(field_flags.cnt_D) AS count_D, SUM(field_flags.cnt_E) AS count_E, SUM(field_flags.cnt_F) AS count_F, SUM(field_flags.cnt_G) AS count_G, SUM(field_flags.cnt_H) AS count_H; -- 格式化成你想要的键值对形式 formatted_output = FOREACH total_counts GENERATE CONCAT('A ', (chararray)count_A), CONCAT('B ', (chararray)count_B), CONCAT('C ', (chararray)count_C), CONCAT('D ', (chararray)count_D), CONCAT('E ', (chararray)count_E), CONCAT('F ', (chararray)count_F), CONCAT('G ', (chararray)count_G), CONCAT('H ', (chararray)count_H); DUMP formatted_output;
这样输出会是一行包含所有字段的计数,和你期望的格式一致。
三、2亿行数据集的性能优化建议
处理这么大的数据集,容易出现内存溢出或数据倾斜,给你几个小技巧:
- 提前设置Reducer数量:
SET mapreduce.job.reduces 30;(根据你的集群资源调整,20-50之间都可以) - 先拿小样本测试逻辑:用
LIMIT dataset 1000;取1000行数据验证过滤和统计逻辑是否正确,再跑全量数据 - 检查字段索引是否正确:你用
$43到$50提取字段,要确认Pig的字段索引是从0开始的,别搞错了字段位置导致提取错误
最后排查小细节
如果还是有异常,一定要去看Pig的任务日志(比如YARN的应用日志),里面会有具体的错误信息——比如是不是数据类型不匹配,或者某个字段的格式有问题导致过滤失败。
内容的提问来源于stack exchange,提问作者Inban
相关产品推荐
相关产品推荐

