Scala环境下如何在Spark UDF中获取DataFrame列值实现跨表查询
问题原因
你遇到的报错是因为col("someId")返回的是Spark SQL的Column类型对象,代表数据集中的一列逻辑引用,不是具体的行级字符串值,无法直接传入要求String类型入参的普通Scala方法lookup中。而且你写的lookup方法是本地执行的算子,每次调用都会触发一次Spark作业,直接在withColumn中逐行调用会产生严重的性能问题,完全不适合分布式计算场景。
最优解决方案:使用等值Join代替自定义lookup
两张DataFrame基于id关联是Spark原生支持的标准操作,性能远高于逐行查询,你可以直接写如下代码:
// 这里用left join是为了保留dataDf中没有匹配到id的行,如果你只需要匹配上的行可以用inner join val resultDf = dataDf.join( lookupdf.select("id", "value"), // 只取需要的关联字段和值字段,减少数据传输 dataDf("someId") === lookupdf("id"), "left" ).withColumnRenamed("value", "lookupVal") // 如果不需要原lookupdf的id列可以后续drop掉 // .drop(lookupdf("id"))
可选方案:使用广播变量 + UDF实现自定义lookup
如果你的lookup逻辑有特殊处理无法直接用join实现,可以把小体量的lookupdf转成广播的Map,再注册成UDF使用:
// 1. 先把lookupdf收集成本地Map,广播到所有 Executor 节点 val lookupMap = lookupdf.select("id", "value").as[(String, String)].collect().toMap val broadcastLookup = spark.sparkContext.broadcast(lookupMap) // 2. 注册自定义UDF val lookupUdf = udf((id: String) => { broadcastLookup.value.getOrElse(id, "") // 匹配不到可以返回默认值,比如空字符串 }) // 3. 在withColumn中调用UDF val resultDf = dataDf.withColumn("lookupVal", lookupUdf(col("someId")))
注意:该方案仅适合lookupdf体量较小的场景,如果lookupdf数据量很大,还是建议使用第一种Join方案。
内容的提问来源于stack exchange,提问作者AOwens
相关产品推荐
相关产品推荐

