ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Hive行转列/列转行实战:从collect_list到explode的底层原理与踩坑记录

Hive行转列/列转行实战:从collect_list到explode的底层原理与踩坑记录 1. 为什么大家都在问“行转列/列转行”做数据处理的人不管是写SQL还是跑Hive任务十有八九都会撞上这么一个问题手里明明是一张规规矩矩的表但业务方要的却是一行里面塞好几个值或者反过来一行拆成多行。这就是典型的“行转列”和“列转行”。说白了行转列就是把原本多行、同类的数据合并成一行里的多个字段比如把每个人的多个手机号拼在一格里或者把多个科目成绩展开成一行的多个列。列转行则恰好相反把一行里的多个字段拆成多行记录比如把“数学、语文、英语”三列拆成三行每行对应一个科目。Hive里做这两个操作核心就是collect_set/collect_list配合concat_ws以及explode配合lateral view。这套东西听起来不复杂但实际用起来坑一个接一个。数据量一大collect出来的顺序根本不受控字符串里带个分隔符一拼接就错位explode碰到空数组直接少一行下游统计就对不上。这篇文章不绕弯子直接讲清楚Hive里行转列/列转行的底层原理、常用写法、踩坑记录和排查手段。适合刚接触Hive的初级工程师也适合被这两种操作折磨过的老手看完至少能少加两天班。2. Hive行转列的核心玩法2.1 从聚合函数说起collect_list和collect_setHive里行转列最常用的工具就是聚合函数collect_list和collect_set。这俩和我们熟悉的sum、max一样都是聚合函数但它们返回的不是单值而是一个数组。collect_list把分组内指定列的所有值收集成一个列表不去重保留重复项。collect_set同样收集成列表但会去重结果里每个值只出现一次。举个例子有一张用户手机号表-- 表数据user_phones -- user_id | phone -- 101 | 13800000001 -- 101 | 13800000002 -- 102 | 13900000001 SELECT user_id, collect_list(phone) AS phone_list FROM user_phones GROUP BY user_id;结果长这样101 | [13800000001,13800000002] 102 | [13900000001]如果换成collect_set遇到重复手机会自动去重。实际业务里到底用哪个取决于需求。比如统计用户购买过的商品标签一般用collect_set因为重复的标签没意义但如果是把每笔交易明细的金额收集起来必须用collect_list保留每一笔否则数据量就对不上。那这个数组生成之后怎么变成“看起来像列”的字符串这就需要concat_ws出场了。concat_ws(separator, array)可以把数组里的元素用指定分隔符拼接起来。比如SELECT user_id, concat_ws(,, collect_list(phone)) AS phones_str FROM user_phones GROUP BY user_id;结果就是101 | 13800000001,13800000002 102 | 13900000001这一步做完多行数据已经变成了“一行多值”的字符串形式。很多报表系统、标签系统里的“多值字段”就是靠这个逻辑存进Hive的。2.2 多列转多行时的group by坑这里有个特别容易踩的坑就是group by的范围没想明白。假如我们要按“用户日期”统计所有商品ID然后拼成一个字符串大多数人会这么写SELECT user_id, stat_date, concat_ws(,, collect_list(product_id)) AS products FROM user_orders GROUP BY user_id, stat_date;这没问题因为需求就是“每个用户每天的商品列表”。但有时候用户需要的是“每个用户所有日期的商品列表”但还希望能看到每个日期对应的商品个数。有人会顺手把stat_date也写进select却忘了加进group by结果报错。Hive里虽然不像传统数据库那么严格开启hive.groupby.orderby.position.alias一些配置后SQL行为会很不一样但对聚合列的要求是一样的——非聚合列必须出现在group by里。正确姿势是SELECT user_id, stat_date, size(collect_list(product_id)) AS cnt, concat_ws(,, collect_list(product_id)) AS products FROM user_orders GROUP BY user_id, stat_date;先分组再收集再拼接顺序别乱。如果真想看“用户所有日期”的数据那就不要stat_date这个粒度只按user_id分组。很多人一开始分不清“粒度”和“列”的关系导致明明想转列结果出来的数据粒度不对下游算指标全错。这里一句话记住group by决定了你要把谁压成一列select里除了聚合函数其他字段必须和group by一致。2.3 动态列转行多行变一行多列刚才说的都是把多行拼成一个字段但有些场景比如做数据透视想把不同分类的值展开成不同的列。Hive里没有MySQL那种pivot函数但可以用case whenmax/sum来模拟。举个常见例子学生成绩表 student_scoresstudent_id | subject | score 1 | math | 90 1 | english | 85 1 | chinese | 88 2 | math | 76 2 | english | 92 2 | chinese | 81现在要转成每个学生一行数学、英语、语文各占一列。SQL这么写SELECT student_id, max(CASE WHEN subject math THEN score END) AS math_score, max(CASE WHEN subject english THEN score END) AS english_score, max(CASE WHEN subject chinese THEN score END) AS chinese_score FROM student_scores GROUP BY student_id;这里能用max也可以用sum因为每个学生科目只有一条记录聚合后非空的只有一个值。为什么不用ifcase when和if能实现同样效果但case when更容易扩展多个条件维护起来清晰。这种方式的缺点是列是写死的。如果你要展开的科目有30个SQL就写成30个case when非常痛苦。实际生产里动态列转行一般分两步先查出来所有枚举值再拼接成SQL或者直接在前端报表工具里做透视不硬写SQL。Hive要真动态展开列就得用map或者str_to_map把数据先变成map再提取键值但这也需要提前知道key。总之Hive对动态列的支持很弱能静态展开就静态展开不能就换思路。3. 列转行的姿势和细节3.1 explode函数把一行炸开成多行列转行的核心是explode。explode是Hive里的UDTFUser Defined Table-Generating Function输入一个数组或map输出多行。比如SELECT explode(array(math,english,chinese)) AS subject;输出的结果是三行math、english、chinese。实际场景里我们经常面对的是“用户订单里存了一个商品ID数组”比如order_id | product_ids 1001 | [p001,p002,p003] 1002 | [p004]要展开成每行一个商品直接SELECT order_id, product_id FROM orders LATERAL VIEW explode(product_ids) t AS product_id;lateral view的作用是把explode产生的多行跟原表的列进行关联形成一张虚拟表。如果不加lateral viewexplode的结果没法和其他列一起select出来这是Hive的语法限制。3.2 lateral view的语法结构拆解很多新手觉得lateral view难记其实它的语法就三块SELECT 原表列, 爆炸出来的列 FROM 原表 LATERAL VIEW explode(要炸的列) 虚拟表别名 AS 新列名;拿上面例子拆解t是虚拟表别名可以任意写比如tmp。product_id是爆炸后新生成的一列后面写SQL就用它。如果explode的数组为空这一行就直接没了。比如某订单product_ids是空数组[]展开后不会产生任何行order_id也会在结果里消失。这在统计订单商品数时会导致订单丢失需要特别小心。处理空数组可以用outer关键字写成SELECT order_id, product_id FROM orders LATERAL VIEW OUTER explode(product_ids) t AS product_id;这样空数组的行也会保留product_id就是NULL。类似left join的效果老版本的Hive用lateral view outer新版本建议直接升级引擎否则很多语法支持不完整。3.3 多个数组同时炸注意笛卡尔积有时候一列不够炸要炸两列。比如一张表user_id | hobbies | skills 1 | [reading] | [python,sql]想同时炸开hobbies和skills有人直接写两个explodeSELECT user_id, hobby, skill FROM user_table LATERAL VIEW explode(hobbies) t1 AS hobby LATERAL VIEW explode(skills) t2 AS skill;这样结果是 readingpython、readingsql 两行。如果hobbies有2个skills有3个就会产生6行这是笛卡尔积。很多时候这不是我们要的我们可能只想要“第1个爱好对应第1个技能”那就不能直接用两个explode。正确的姿势是使用posexplode它能同时返回数组的下标和值。先炸第一个数组保留下标再炸第二个数组用下标关联SELECT t1.user_id, t1.hobby, t2.skill FROM ( SELECT user_id, hobby, row_number() OVER (PARTITION BY user_id ORDER BY 1) AS rn FROM user_table LATERAL VIEW posexplode(hobbies) tmp AS pos, hobby ) t1 JOIN ( SELECT user_id, skill, row_number() OVER (PARTITION BY user_id ORDER BY 1) AS rn FROM user_table LATERAL VIEW posexplode(skills) tmp AS pos, skill ) t2 ON t1.user_id t2.user_id AND t1.rn t2.rn;这个写法稍微绕一点但能精准控制关联逻辑。实际业务里“两个数组长度不一致”的情况特别常见用笛卡尔积会爆炸用下标关联才能对得上位次。3.4 列转行后与原表其他字段做关联的常见坑上面已经提到lateral view是跟在每张表后面的。如果原表本身就是个子查询那lateral view应该放在子查询的最外层不要在子查询内部炸完再关联那样性能很差还会造成不必要的数据膨胀。举个例子有人这么写SELECT t.user_id, s.product_id FROM ( SELECT user_id, concat_ws(,, collect_list(product_id)) AS products FROM user_orders GROUP BY user_id ) t LATERAL VIEW explode(split(t.products, ,)) s AS product_id;先聚合拼接再炸开这不是脱裤子放屁吗如果只是为了展开直接炸原数组就好了。但有一种情况例外如果你要的是“先分组再展开去重后的结果”那可以collect_set之后再炸。不过多数情况下既然能直接explode就不要再绕一圈只会增加额外的序列化和反序列化开销。4. 实战场景从订单明细到标签宽表4.1 需求描述与数据准备举一个实际生产的例子。假设有订单明细表order_detailorder_id | user_id | product_id | category_id | sale_amount 1001 | 201 | p001 | c01 | 99.00 1002 | 201 | p002 | c02 | 159.00 1003 | 202 | p003 | c01 | 79.00 1004 | 202 | p004 | c03 | 299.00 1005 | 203 | p005 | c02 | 49.00需求是生成一张用户维度宽表每个用户一行包含用户ID购买过的所有商品ID用逗号拼接购买过的所有类目ID去重后用竖线拼接每个类目的销售额合计订单数量这就是一个典型的“多行转一行多列”场景。大部分字段可以通过普通的聚合完成但“商品ID拼接”“类目ID拼接”就得靠collect_list/collect_set。4.2 先聚合后拼接还是先拼接后聚合这里再强调一次SQL里是先group by聚合再对聚合结果做展示。所有collect函数都是聚合函数所以它必然在group by的语境下出现。上面的需求可以一行SQL完成SELECT user_id, concat_ws(,, collect_list(product_id)) AS products, concat_ws(|, collect_set(category_id)) AS categories, SUM(sale_amount) AS total_amount, COUNT(*) AS order_cnt FROM order_detail GROUP BY user_id;结果201 | p001,p002 | c01|c02 | 258.00 | 2 202 | p003,p004 | c01|c03 | 378.00 | 2 203 | p005 | c02 | 49.00 | 1这个SQL看起来简单但细节都藏在函数选择上collect_list(product_id)保留所有商品即使重复也保留因为用户可能多次购买同一商品。collect_set(category_id)去重因为类目只要出现过一次就够了。concat_ws的拼接顺序不保证Hive里collect_list的顺序依赖map端输出和reduce端的处理不能当作有序列表用。如果对顺序有要求需要用sort_array包一层或者提前用row_number排序。4.3 遇到顺序要求怎么破业务里还真有“产品列表按购买时间倒序”的需求。这个在Hive里做起来比较绕但有一个办法先给数据排序并带上序号再用collect_list收集按序号排序的对象。比如SELECT user_id, concat_ws(,, collect_list(product_id)) AS products FROM ( SELECT user_id, product_id, row_number() OVER (PARTITION BY user_id ORDER BY order_time DESC) AS rn FROM order_detail ) t GROUP BY user_id;这么说吧内部子查询已经按时间倒序排好了但collect_list收集时依然不保证顺序。真正的解法是把行号也一起收集在select之后对数组进行二次排序SELECT user_id, concat_ws(,, sort_array(collect_list(concat_ws(:, lpad(rn, 10, 0), product_id)))) AS sorted_products FROM ( SELECT user_id, product_id, row_number() OVER (PARTITION BY user_id ORDER BY order_time DESC) AS rn FROM order_detail ) t GROUP BY user_id;这里把rn和product_id拼接成“0000000001:p001”这种字符串sort_array按字典序排完再拆出来。虽然绕但稳定。如果不这么干你会发现两次跑任务拼出来的商品顺序可能不一样这数据直接上线会被业务方骂死。4.4 把宽表再拆回明细用户类目透视数据仓库里经常要反过来把宽表拆回明细做透视统计。比如上面生成的用户类目组合要拆成每个用户类目一行方便用BI工具做交叉筛选。这时候用lateral view explode(split(...))SELECT user_id, category_id FROM user_wide_table LATERAL VIEW explode(split(categories, \\|)) t AS category_id;注意split的分隔符转义。如果分割符是竖线|它是正则表达式里的特殊字符必须写成\\|。要是逗号直接split(categories, ,)就行。这个转义问题是我见过报错最多的点之一。很多人写split(categories, |)出来结果是每个字符都拆开了因为|在正则里表示“或”匹配空字符串结果每个字符之间都被切了。这里给出一个经验规则在Hive SQL里凡是遇到.、|、*、、?、(、)、[、]、{、}、\、^、$这些正则字符需要再套一层转义。最常见的几个-- 竖线 split(str, \\|) -- 点 split(str, \\.) -- 星号 split(str, \\*)如果嫌转义麻烦可以用regexp_replace先把字符替换成逗号再split但性能会稍微差一点。4.5 用str_to_map解析key-value结构的列转行Hive里还有一种情况原始数据是key1:value1;key2:value2这种字符串需要拆成KV对再转成多行多列。比如某个埋点日志里的扩展字段ext_info os:android;channel:huawei;version:1.2.0先用str_to_map(ext_info, ;, :)转成map再explodeSELECT log_id, key_name, key_value FROM log_table LATERAL VIEW explode(str_to_map(ext_info, ;, :)) t AS key_name, key_value;这个写法在解析半结构化数据时非常实用也能用来处理动态列。比如上面的例子最终输出三行os→android、channel→huawei、version→1.2.0。如果后续要转列再用case when根据key_name提取或者直接map的键访问。str_to_map还有几个参数比如可以指定两个分界符还能设置是否允许key重复。实际使用中一定要确认数据里分隔符不会出现在value里否则解析会错位。我之前处理过一批埋点数据value里带了个分号结果整个map就乱了排查半天才定位到是脏数据问题。所以在解析这类KV串之前最好先做一次合法性过滤比如用regexp_replace去掉敏感分隔符。5. 数据量和性能视角下的行转列/列转行5.1 为什么说collect_list在Map端做更高效Hive执行聚合时如果group by的key分布均匀且数据量不大聚合可以在Map端完成也就是所谓的Map端聚合。collect_list这种聚合逻辑也能走Map端好处是Shuffle数据量小Reduce端需要处理的数据少。但有一个问题是collect_list会把所有的值都留在内存里。如果一个分组下有几十万行收集出来的数组可能非常大导致OOM。遇到这种情况你可以提高hive.map.aggr.hash.percentmemory但治标不治本。在收集前先用row_number截断比如每个分组只保留前100条。把大字段拆出来单独处理不要跟其他聚合混在一个任务里。我实际处理过一个用户行为表某个用户一年有50万条行为记录直接collect_list把Reducer内存打爆了。后来改成先按天聚合再按用户汇总分两层才解决。具体过程是第一步按用户天聚合生成当天的行为列表INSERT OVERWRITE TABLE user_behavior_day SELECT user_id, stat_date, concat_ws(,, collect_list(action)) AS actions FROM raw_behavior GROUP BY user_id, stat_date;第二步按用户聚合把天的列表再拼一次INSERT OVERWRITE TABLE user_behavior_total SELECT user_id, concat_ws(,, collect_list(actions)) AS all_actions FROM user_behavior_day GROUP BY user_id;虽然多跑一次任务但是每个分组的大小从50万变到了最多365内存压力小得多。5.2 explode的数据膨胀预判结果行数列转行的代价是数据膨胀。一个数组有10个元素炸完就是10行。如果一个表有1亿行每行平均炸10个最终结果就是10亿行。这种膨胀在下游join时特别容易造成数据倾斜。比如用lateral view explode炸开某个巨长的数组会把同一个key的所有行都派发到同一个Reduce因为要保证同一个key的数据在一起才能join。如果某个key的数组特别长其他key都很短那这个Reduce就会成为瓶颈。优化思路通常是如果能先过滤就先过滤数组里的元素比如把长度超过阈值的截断再炸。如果不能过滤就考虑把长数组单独抽出来处理或者用map join打散。避免在explode之后再shuffle join尽量用map join但需要把小表加载到内存里。5.3 数据倾斜的表现和简查方法做行转列/列转行时最典型的数据倾斜症状是跑任务卡在某个Reduce进度一直是99%其他Reduce早就完了。如果你用Tez引擎打开YARN的ResourceManager界面看一下某个Container的内存或CPU明显比其他的高基本就是倾斜。排查步骤看SQL里有没有group by或join倾斜大概率发生在这两个算子。最直接的方法是查分组分布SELECT user_id, COUNT(*) AS cnt FROM order_detail GROUP BY user_id ORDER BY cnt DESC LIMIT 10;如果某个用户的记录数比其他用户高出几个数量级那就是倾斜根因。对症下药聚合倾斜开启hive.groupby.skewindatatrue强制拆成两轮聚合第一轮用随机前缀打散第二轮去掉前缀聚合。这个参数我已经用了很多次对简单聚合非常有效。join倾斜先把大key筛选出来单独处理或者用map join把大key广播。collect_list倾斜上面的“二楼聚合”思路。6. 常见报错与问题排查速查表下面把Hive里做行转列、列转行最常见的报错和现象整理成一个表方便遇到问题直接对号入座。报错信息或现象可能原因解决思路FAILED: SemanticException [Error 10081]: UDTFs are not supported outside the SELECT clauseexplode用在了select之外的位置比如where中explode只能在select或lateral view中使用Invalid table alias or column referencelateral view的虚拟表别名或列名写错检查lateral view explode(col) 别名 AS 列名的写法确保列名唯一结果里每个字符都被拆开split用了正则特殊字符比如|没转义把分隔符改成\|或\\.collect_list出来的顺序不稳定Hive不保证聚合结果顺序用sort_array包一层或提前拼接行号再排序空数组或null导致行丢失lateral view会忽略空结果改用lateral view outer explodejava.lang.NoClassDefFoundError: org/apache/hadoop/crypto运行时环境缺Hadoop加密库或依赖冲突检查Hive lib下是否有hadoop-crypto相关jar检查HADOOP_CLASSPATH重启hive服务GROUP BY结果比预期少很多行某字段被误写进group by或collect_set自动去重明确需求是去重还是不去重去重用set不去重用list内存溢出OOMcollect_list收集的数据量过大分两层聚合先细粒度再粗粒度聚合结果字段变成了NULLcase when没有匹配条件且没有用max/sum包住聚合函数要包住case when否则每行一个值取出来就变成NULL了SQL语法在SparkSQL里能跑Hive里报错引擎差异Hive对UDTF和lateral view支持更严格以Hive官方文档为准或者换基于Hive的底层能力来写这里单独说一下java.lang.NoClassDefFoundError。有些刚配好的Hive环境跑行转列时突然报这个其实和行转列本身没关系纯粹是环境依赖缺失。常见原因有安装Hive时只拷贝了hive的lib没拷贝Hadoop的common和crypto jar。Hadoop版本和Hive版本不匹配导致某些类找不到。提交任务时classpath没带全。解决办法很简单找到hadoop-crypto-*.jar并拷到Hive的lib目录下或者设置HADOOP_CLASSPATH指向Hadoop的share目录。如果改完还报错记得重启Hive服务像HS2这种长时间运行的服务不加HIVE_CLI_SERVICE_LOCATION这些配置的话类路径经常不刷新。7. 从Hive到其他引擎语法差异提醒很多人平时用Hive偶尔也去写SparkSQL、Flink SQL会发现行转列/列转行的用法基本相通但细节上还是有差异。SparkSQL里collect_list/collect_set和concat_ws的组合完全一样explode和lateral view也兼容。但有几个不同点SparkSQL里split的正则转义和Hive一样需要双重转义。SparkSQL的posexplode语法完全一样。但SparkSQL对null的处理有时更宽松比如concat_ws在Hive里会自动忽略nullSparkSQL在部分版本里会把null转成空字符串这一点对接下游数据时要注意。Flink SQL里lateral view换成了cross join unnest()虽然理解逻辑一致但函数名和写法完全不同。如果是从Hive迁到Flink不能直接复制SQL要改语法。所以我建议团队里如果同时用多个引擎最好封装一层标准的UDF或视图把“行转列/列转行”的逻辑固定下来。比如统一要求所有拼接字段必须用concat_ws所有拆分必须用自定义的split_explodeUDTF这样就可以屏蔽引擎差异。8. 一个自定义UDF替代方案如果你被split转义搞得头痛或者需要更复杂的行转列逻辑可以写一个自定义UDF。我提供一个最简单但实用的思路写一个UDF把任意字符串按任意分隔符拆成数组再配合Hive自带的explode使用。这个UDF的好处是可以自定义处理空白、过滤空串避免脏数据导致下游多出空行。伪代码示意public class SplitAndFilter extends UDF { public ListString evaluate(String str, String sep) { if (str null || str.isEmpty()) return new ArrayListString(); String[] parts str.split(sep); ListString result new ArrayListString(); for (String part : parts) { if (part ! null !part.trim().isEmpty()) { result.add(part.trim()); } } return result; } }打包成jar后在Hive里ADD JAR /path/to/split_filter.jar; CREATE TEMPORARY FUNCTION split_filter AS com.example.SplitAndFilter; SELECT user_id, item FROM user_table LATERAL VIEW explode(split_filter(categories, \\|)) t AS item;其实这个功能用splitexplode也能实现但自定义UDF可以把过滤逻辑写进去比如去掉空串、去空格、统一大小写。在数据质量和逻辑复用要求高的场景值得投入一点开发时间。9. 实操心得与最后一点建议做了这么多年数据开发行转列/列转列几乎每周都要用以下几条经验算是久病成医分享出来给大家。第一写Hive SQL前先想清楚粒度和数据形态。你不会希望在写完SQL之后才发现group by少了一个维度导致结果多出很多行也不会希望明明做了去重下游却发现重复。我的习惯是把目标表的字段一个个记下来标注“是否聚合字段”“聚合函数是什么”“分组字段是什么”再动手写SQL。第二优先级排序能不用正则就不正则能不用字符串拼接就不用字符串拼接。数组和map在Hive里是一等公民直接用它们做中间结果最后在出口处再转成字符串。比如先collect_list成数组再用array_contains做过滤比先把数组转成字符串再用like去匹配性能高一个数量级。第三不要忽略大小写和字段类型。collect_set对字符串和数值的区分很严格同为数字的1和字符串1会被当成两个值。所以ETL里要尽早统一字段类型否则行转列后出现“看起来重复但实际上没重复”的奇葩数据。第四遇到复杂需求先拆步。行转列后要做多个汇总指标时不要硬塞一个超大SQL分步骤生成中间表每步验证一次结果。调试成本远比运行成本高尤其我见过不少人直接在生产环境跑查询一旦超时或者OOM反而耽误更多时间。最后再分享一个我常用的调试技巧在开发环境把数据量限制到几百条把所有结果打印到本地肉眼核对一遍中间结果。不要急着直接上全量数据因为你根本看不出数据是否正确。很多所谓“行转列结果不对”的问题其实都是开发阶段没有仔细检查中间结果要么是源数据重复要么是join出了一对多跟函数本身的逻辑没有关系。Hive的行转列/列转行虽然看起来就是几个函数拆来拆去但真正吃透之后你会发现无论数据形态怎么变其实核心思想都是“以什么粒度聚合、用什么函数展开”。把这个底层逻辑想通了遇到任何新函数也只是换个工具而已。
返回列表