在Spark的groupByKey和countByKey中使用JodaTime

Joe*_*Joe 5 jodatime apache-spark

我有一个非常简单的Spark程序(在Clojure中使用Flambo,但应该很容易理解).这些是JVM上的所有对象.我正在测试一个local实例(虽然我猜测Spark仍然会序列化和反序列化).

(let [dt (t/date-time 2014)
      input (f/parallelize sc [{:the-date dt :x "A"}
                               {:the-date dt :x "B"}
                               {:the-date dt :x "C"}
                               {:the-date dt :x "D"}])
      by-date (f/map input (f/fn [{the-date :the-date x :x}] [the-date x])))
Run Code Online (Sandbox Code Playgroud)

输入是四个元组的RDD,每个元组具有相同的日期对象.第一个映射产生date => x的键值RDD.

内容input如预期:

=> (f/foreach input prn)
[#<DateTime 2014-01-01T00:00:00.000Z> "A"]
[#<DateTime 2014-01-01T00:00:00.000Z> "B"]
[#<DateTime 2014-01-01T00:00:00.000Z> "C"]
[#<DateTime 2014-01-01T00:00:00.000Z> "D"]
Run Code Online (Sandbox Code Playgroud)

为了清楚,平等和.hashCode工作日期对象:

=> (= dt dt)
true
=> (.hashCode dt)
1260848926
=> (.hashCode dt)
1260848926
Run Code Online (Sandbox Code Playgroud)

它们是JodaTime的DateTime实例,它按预期实现等于.

当我尝试时countByKey,我得到了预期的:

=> (f/count-by-key by-date)
{#<DateTime 2014-01-01T00:00:00.000Z> 4}
Run Code Online (Sandbox Code Playgroud)

但是,当我groupByKey,它似乎不起作用.

=> (f/foreach (f/group-by-key by-date) prn)
[#<DateTime 2014-01-01T00:00:00.000Z> ["A"]]
[#<DateTime 2014-01-01T00:00:00.000Z> ["B"]]
[#<DateTime 2014-01-01T00:00:00.000Z> ["C"]]
[#<DateTime 2014-01-01T00:00:00.000Z> ["D"]]
Run Code Online (Sandbox Code Playgroud)

密钥都是相同的,所以我希望结果是一个条目,以日期为关键字和["A", "B", "C", "D"]值.发生了一些事情,因为值都是列表.

不知何故groupByKey,没有正确地将键等同起来.但是countByKey.这两者有什么区别?我怎样才能让它们表现得一样?

有任何想法吗?

Joe*_*Joe 3

我越来越接近答案了。我认为这属于答案部分而不是问题部分。

这按键分组,变成本地收集,提取第一项(日期)。

=> (def result-dates (map first (f/collect (f/group-by-key by-date))))
=> result-dates
(#<DateTime 2014-01-01T00:00:00.000Z>
 #<DateTime 2014-01-01T00:00:00.000Z>
 #<DateTime 2014-01-01T00:00:00.000Z>
 #<DateTime 2014-01-01T00:00:00.000Z>)
Run Code Online (Sandbox Code Playgroud)

哈希码都是相同的

=> (map #(.hashCode %) result-dates)
(1260848926
 1260848926
 1260848926 
 1260848926)
Run Code Online (Sandbox Code Playgroud)

毫秒都是相同的:

=> (map #(.getMillis %) result-dates)
(1388534400000
 1388534400000
 1388534400000
 1388534400000)
Run Code Online (Sandbox Code Playgroud)

equals失败了,但是isEquals成功了

=> (.isEqual (first result-dates) (second result-dates))
true

=> (.equals (first result-dates) (second result-dates))
false
Run Code Online (Sandbox Code Playgroud)

文档.equals

根据毫秒时刻和时间顺序比较此对象与指定对象的相等性

它们的毫秒数都是相等的,它们的年表似乎是:

=> (map #(.getChronology %) result-dates)
(#<ISOChronology ISOChronology[UTC]>
 #<ISOChronology ISOChronology[UTC]>
 #<ISOChronology ISOChronology[UTC]>
 #<ISOChronology ISOChronology[UTC]>)
Run Code Online (Sandbox Code Playgroud)

然而,年表并不等同。

=> (def a (first result-dates))
=> (def b (second result-dates))

=> (= (.getChronology a) (.getChronology b))
false
Run Code Online (Sandbox Code Playgroud)

虽然哈希码确实

=> (= (.hashCode (.getChronology a)) (.hashCode (.getChronology b)))
true
Run Code Online (Sandbox Code Playgroud)

但是joda.time.Chronology没有提供自己的equals方法,而是继承自Object,它只使用引用相等。

我的理论是,这些日期都用它们自己的、不同的、构造的 Chronology 对象进行反序列化,但 JodaTime 有自己的序列化器,可能可以处理这个问题。也许自定义Kryo序列化器会在这方面有所帮助。

目前,我在 Spark 中使用 JodaTime 的解决方案是通过调用, 或 a而不是org.joda.time.DateTime来使用org.joda.time.InstanttoInstantjava.util.Date

两者都涉及丢弃时区信息,这并不理想,因此如果有人有更多信息,将非常受欢迎!