【发布时间】:2017-11-13 14:13:59
【问题描述】:
我在 spark 中有一个 UDF(在 EMR 上运行),它是用 scala 编写的,它使用用于 scala 的 uaparser 库(uap-scala)从用户代理解析设备。在小型集合上工作时它工作正常(5000 行),但在大型集合(2M)上运行时它工作得非常慢。 我尝试收集 Dataframe 以列出并在驱动程序上循环它,这也很慢,是什么让我相信 UDF 在驱动程序而不是工人上运行
- 我怎样才能确定这一点?有人有其他理论吗?
- 如果是这样,为什么会发生这种情况?
这是 udf 代码:
def calcDevice(userAgent: String): String = {
val userAgentVal = Option(userAgent).getOrElse("")
Parser.get.parse(userAgentVal).device.family
}
val calcDeviceValUDF: UserDefinedFunction = udf(calcDevice _)
用法:
.withColumn("agentDevice", udfDefinitions.calcDeviceValUDF($"userAgent"))
谢谢 尼尔
【问题讨论】:
-
userAgent的某些值可以被不同的行共享吗? -
是的,用户代理会重复自己,但列表不是很小
-
解析可能很昂贵 - 如果使用缓存,命中率会是多少?
标签: scala apache-spark emr ua-parser