如何比较两个数据集?

Hem*_*nth 12 scala apache-spark fastutil

我正在运行一个spark应用程序,它从几个hive表(IP地址)读取数据,并将数据集中的每个元素(IP地址)与来自其他数据集的所有其他元素(IP地址)进行比较.最终结果将是这样的:

+---------------+--------+---------------+---------------+---------+----------+--------+----------+
|     ip_address|dataset1|dataset2       |dataset3       |dataset4 |dataset5  |dataset6|      date|
+---------------+--------+---------------+---------------+---------+----------+--------+----------+
| xx.xx.xx.xx.xx|     1  |              1|              0|        0|         0|      0 |2017-11-06|
| xx.xx.xx.xx.xx|     0  |              0|              1|        0|         0|      1 |2017-11-06|
| xx.xx.xx.xx.xx|     1  |              0|              0|        0|         0|      1 |2017-11-06|
| xx.xx.xx.xx.xx|     0  |              0|              1|        0|         0|      1 |2017-11-06|
| xx.xx.xx.xx.xx|     1  |              1|              0|        1|         0|      0 |2017-11-06|
---------------------------------------------------------------------------------------------------
Run Code Online (Sandbox Code Playgroud)

为了进行比较,我将语句的dataframes结果转换hiveContext.sql("query")为Fastutil对象.像这样:

val df= hiveContext.sql("query")
val dfBuffer = new it.unimi.dsi.fastutil.objects.ObjectArrayList[String](df.map(r => r(0).toString).collect())
Run Code Online (Sandbox Code Playgroud)

然后,我使用一个iterator迭代每个集合并使用FileWriter.将行写入文件.

val dfIterator = dfBuffer.iterator()
while (dfIterator.hasNext){
     val p = dfIterator.next().toString
     //logic
}
Run Code Online (Sandbox Code Playgroud)

我正在运行应用程序 --num-executors 20 --executor-memory 16g --executor-cores 5 --driver-memory 20g

该过程总共运行大约18-19小时,大约4-5百万条记录,每天进行一对一的比较.

然而,当我检查了应用主界面,我注意到,没有活动发生的初始转换后,dataframes到fastutil collection objects完成(这需要工作启动仅几分钟后).我看到代码中使用的count和collect语句产生新的作业,直到转换完成.之后,比较运行时不会启动任何新作业.

  • 这意味着什么?这是否意味着分布式处理根本没有发生?

  • 我知道收集对象不被视为RDD,
    这可能是原因吗?

  • 如果不使用分配的资源,spark如何执行我的程序?

任何帮助将不胜感激,谢谢!

Jac*_*ski 8

行后:

val dfBuffer = new it.unimi.dsi.fastutil.objects.ObjectArrayList[String](df.map(r => r(0).toString).collect())
Run Code Online (Sandbox Code Playgroud)

ESP.上述部分内容:

df.map(r => r(0).toString).collect()
Run Code Online (Sandbox Code Playgroud)

这collect是最重要的事情,没有执行任何Spark作业dfBuffer(这是一个常规的本地JVM数据结构).

这是否意味着分布式处理根本没有发生?

正确.collect将所有数据带到驱动程序运行的单个JVM上(这正是你不应该这样做的原因,除非......你知道你在做什么以及它可能导致什么问题).

我认为以上回答了所有其他问题.


比较两个数据集(以Spark和分布式方式)的问题的可能解决方案是join使用参考数据集的数据集,并count比较记录数是否没有变化.