【发布时间】:2016-08-31 00:39:34
【问题描述】:
我有两个结构如下的文件
文件 1
gnk_id, matchId, timestamp
文件 2
gnk_matchid, matchid
如果file1.gnk_id = file2.gnk_machid,我想用文件2中matchid的值更新文件1中gnk_id的值。
为此,我在 Spark 中创建了两个数据框。我想知道我们是否可以更新 Spark 中的值?如果没有,是否有任何解决方法可以提供更新的最终文件?
更新
我做了这样的事情
case class GnkMatchId(gnk: String, gnk_matchid: String)
case class MatchGroup(gnkid: String, matchid: String, ts: String)
val gnkmatchidRDD = sc.textFile("000000000001").map(_.split(',')).map(x => (x(0),x(1)) )
val gnkmatchidDF = gnkmatchidRDD.map( x => GnkMatchId(x._1,x._2) ).toDF()
val matchGroupMr = sc.textFile("part-00000").map(_.split(',')).map(x => (x(0),x(1),x(2)) ).map( f => MatchGroup(f._1,f._2,f._3.toString) ).toDF()
val matchgrp_joinDF = matchGroupMr.join(gnkmatchidDF,matchGroupMr("gnkid") === gnkmatchidDF("gnk_matchid"),"left_outer")
matchgrp_joinDF.map(x => if(x.getAs[String]("gnk_matchid").length != 0 ) {MatchGroup(x.getAs[String]("gnk_matchid"), x.getAs[String]("matchid"),x.getAs[String]("ts"))} else {MatchGroup(x.getAs[String]("gnkid"), x.getAs[String]("matchid"),x.getAs[String]("ts"))}).toDF().show()
但在最后一步,NULLpointerEXception 失败了
【问题讨论】:
-
在这里回答了类似的问题stackoverflow.com/questions/36800174/…
标签: apache-spark apache-spark-sql