【发布时间】:2019-02-12 18:43:56
【问题描述】:
我是 Zeppelin 的新手,也许我的问题很幼稚。一开始,我得到的基本数据是这样的:
import org.apache.spark.sql.functions.sql
val dfOriginal = sql("SELECT CAST(event_type_id AS STRING), event_time FROM sl_event SORT BY event_time LIMIT 200")
+-------------+--------------------+
|event_type_id| event_time|
+-------------+--------------------+
| 23882|2018-05-03 11:41:...|
| 23882|2018-05-03 11:41:...|
| 23882|2018-05-03 11:41:...|
| 25681|2018-05-03 11:41:...|
| 23882|2018-05-03 11:41:...|
| 2370|2018-05-03 11:41:...|
| 23882|2018-05-03 11:41:...|
...
我有 200 条这样的记录。
我计算偶数类型的出现次数如下:
val dfIndividual = dfOriginal.groupBy("event_type_id").count().sort(-col("count"))
dfIndividual.show(200)
我很困惑:每当我执行这个(在 Zeppelin 中)时,我都会得到不同的结果。例如:
+-------------+-----+
|event_type_id|count|
+-------------+-----+
| 24222| 30|
| 10644| 16|
| 21164| 9|
...
或 - 几秒钟后:
+-------------+-----+
|event_type_id|count|
+-------------+-----+
| 5715| 34|
| 3637| 19|
| 3665| 17|
| 9280| 13|
...
这两个结果之间的差异让我非常害怕。问题出在哪里?是齐柏林飞艇吗?基础火花?如何保证我会在这里得到可重现的结果?
【问题讨论】:
标签: scala apache-spark apache-spark-sql apache-zeppelin