Scala中Spark行内数组扁平化问题求助
解决Spark Scala中Row嵌套序列的扁平化问题
我完全懂你的困扰——手动罗列几十列不仅麻烦还容易出错,谁都不想干这种重复活。你遇到的问题核心有两个:一是之前的代码把30 to 33范围内的字段打包成了Vector作为Row的单个嵌套元素,不符合你要的扁平字段需求;二是用flatMap报错,因为它的作用是把单个输入元素拆分成多个输出元素(比如把一个Row拆成多个Row),但你要的是修改单个Row的结构,让它的字段更扁平,这根本不是flatMap的适用场景。
正确解决方案:用序列拼接构造扁平Row
我们可以通过拼接序列的方式,把所有需要的字段整合成一个平级序列,再用Row.fromSeq生成目标Row。不管你有多少列,只需要调整字段范围即可,完全不用手动逐个罗列:
myfile.filter(!_.contains("BatchNumber")) .map(line => line.split("\t")) .map { line => // 拼接所有需要的字段:第一个字段 + 30到33的元素 + 最后一个字段 val flatFields = Seq(line(0)) ++ (30 to 33).map(line(_)) ++ Seq(line(1)) Row.fromSeq(flatFields) } .take(2)
为什么这个方法有效?
Seq(line(0))把第一个字段包装成序列,方便后续和其他序列拼接;(30 to 33).map(line(_))自动生成30到33索引对应的字段序列,每个元素都是独立的;++操作符把三个序列拼接成一个完全平级的大序列;Row.fromSeq(flatFields)直接把平级序列转换成Row,每个序列元素对应Row的一个独立字段。
执行这段代码后,你就能得到预期的结果:
Array([C002,323.000, 0.000, 0.000, 323.000,2RFG], [C002,33.000, 0.000, 0.000, 33.000,2RFG])
灵活扩展:处理非连续字段
如果你的目标字段不是连续范围(比如需要30、32、35-40这样的组合),也能轻松调整序列逻辑,比如:
val flatFields = Seq(line(0), line(30), line(32)) ++ (35 to 40).map(line(_)) ++ Seq(line(1))
同样能生成符合需求的扁平Row,完全不用手写几十个line(i)。
内容的提问来源于stack exchange,提问作者Cauchy Wu
相关产品推荐
相关产品推荐

