【发布时间】:2017-06-08 20:26:16
【问题描述】:
解决以下用例的优化或最佳性能方法是什么
考虑一个包含 100 万行和 100 列的数据框,我们对其中的 1 列感兴趣 - 消息。我需要根据消息中匹配关键字的存在条件构建 3 个新列。
- 消息:堆栈溢出对代码开发的贡献是 一天比一天增加
- flag1 关键字:堆栈、松弛
- flag2 关键字:twitter、facebook、whatsapp
- flag3 关键字:流、运行、增加
预期输出:(message,flag1,flag2,flag3) 堆栈溢出对代码开发的贡献日益增加,1,0,0
方法 1
val tempDF = df.withColumn("flag1",computeFlag(col("message"))).withColumn("flag2",computeFlag(col("message"))).withColumn("flag3",computeFlag(col("message")))
方法 2
val tempDF = df.withColumn("flagValues",computeMultipleFlags(col("message"))).withColumn("_tmp", split($"flagValues","#")).select($"message",$"_tmp".getItem(0).as("flag1"),$"_tmp".getItem(1).as("commercial"),$"_tmp".getItem(2).as("flag2"),$"_tmp".getItem(3).as("flag3")).drop("_tmp")
UDF : computeFlag 根据各个关键字列表的精确匹配返回 1 或 0
UDF : computeMultipleFlags 根据 flag1、flag 2 和 flag 3 各自关键字的精确匹配返回 # 分隔结果 1 或 0:示例 1#0#0
我已经使用这两种方法解决了,但是看到/感觉方法 2 表现更好。请指教。
Spark 数据帧默认是并行化的,但是这种情况如何 采用方法 1。flag1、flag2、flag3 列将在 并行还是顺序?
Spark 数据框会自动并行处理我的输入列吗 “消息”:针对列的多行多线程
计算?
【问题讨论】:
标签: scala apache-spark apache-spark-sql user-defined-functions