【问题标题】:Convert from clojure.lang.LazySeq to type org.apache.spark.api.java.JavaRDD从 clojure.lang.LazySeq 转换为 org.apache.spark.api.java.JavaRDD 类型
【发布时间】:2015-08-25 14:49:06
【问题描述】:

我在 clojure 中开发了一个函数来从最后一个非空值填充一个空列,我假设这是可行的,给定

(:require [flambo.api :as f])

(defn replicate-val
  [ rdd input ]
  (let [{:keys [ col ]} input
    result (reductions (fn [a b]
                         (if (empty? (nth b col))
                           (assoc b col (nth a col))
                           b)) rdd )]
(println "Result type is: "(type result))))

知道了:

;=> "Result type is:  clojure.lang.LazySeq"

问题是如何使用 flambo(火花包装器)将其转换回 JavaRDD 类型

我在let表单中尝试(f/map result #(.toJavaRDD %))尝试转换为JavaRDD类型

我收到了这个错误

"No matching method found: map for class clojure.lang.LazySeq"

这是预期的,因为结果的类型是clojure.lang.LazySeq

问题是如何进行这种转换,或者如何重构代码以适应这种情况。

这是一个示例输入 rdd:

(type rdd) ;=> "org.apache.spark.api.java.JavaRDD"

但看起来像:

[["04" "2" "3"] ["04" "" "5"] ["5" "16" ""] ["07" "" "36"] ["07" "" "34"] ["07" "25" "34"]]

需要的输出是:

[["04" "2" "3"] ["04" "2" "5"] ["5" "16" ""] ["07" "16" "36"] ["07" "16" "34"] ["07" "25" "34"]]

谢谢。

【问题讨论】:

    标签: java clojure apache-spark clojure-java-interop flambo


    【解决方案1】:

    首先,RDD 是不可迭代的(不要实现ISeq)所以你不能使用reductions。忽略访问先前记录的整个想法是相当棘手的。首先,您不能直接访问另一个分区中的值。此外,只有不需要改组的转换才能保留顺序。

    这里最简单的方法是使用具有显式顺序的数据框和窗口函数,但据我所知,Flambo 没有实现所需的方法。始终可以使用原始 SQL 或访问 Java/Scala API,但如果您想避免这种情况,可以尝试以下管道。

    首先让我们创建一个广播变量,其中包含每个分区的最后一个值:

    (require '[flambo.broadcast :as bd])
    (import org.apache.spark.TaskContext)
    
    (def last-per-part (f/fn [it]
      (let [context (TaskContext/get) xs (iterator-seq it)]
      [[(.partitionId context) (last xs)]])))
    
    (def last-vals-bd
     (bd/broadcast sc
       (into {} (-> rdd (f/map-partitions last-per-part) (f/collect)))))
    

    接下来是实际工作的一些助手:

    (defn fill-pair [col]
      (fn [x] (let [[a b] x] (if (empty? (nth b col)) (assoc b col (nth a col)) b))))
    
    (def fill-pairs
      (f/fn [it] (let [part-id (.partitionId (TaskContext/get)) ;; Get partion ID
                       xs (iterator-seq it) ;; Convert input to seq
                       prev (if (zero? part-id) ;; Find previous element
                         (first xs) ((bd/value last-vals-bd) part-id))        
                       ;; Create seq of pairs (prev, current)
                       pairs (partition 2 1 (cons prev xs))
                       ;; Same as before
                       {:keys [ col ]} input
                       ;; Prepare mapping function
                       mapper (fill-pair col)]
                   (map mapper pairs))))
    

    最后你可以使用fill-pairsmap-partitions

    (-> rdd (f/map-partitions fill-pairs) (f/collect))
    

    这里隐藏的假设是分区的顺序遵循值的顺序。一般情况下它可能是也可能不是,但如果没有明确的顺序,它可能是你能得到的最好的。

    替代方法是zipWithIndex,交换值的顺序并使用偏移量执行连接。

    (require '[flambo.tuple :as tp])
    
    (def rdd-idx (f/map-to-pair (.zipWithIndex rdd) #(.swap %)))
    
    (def rdd-idx-offset
      (f/map-to-pair rdd-idx
        (fn [t] (let [p (f/untuple t)] (tp/tuple (dec' (first p)) (second p))))))
    
    (f/map (f/values (.rightOuterJoin rdd-idx-offset rdd-idx)) f/untuple)
    

    接下来您可以使用与以前类似的方法进行映射。

    编辑

    快速注释on using atoms。缺乏参考透明度以及您正在利用给定实现的附带属性而不是合同的问题是什么。 map 语义中没有任何内容需要按给定顺序处理元素。如果内部实现发生变化,它可能不再有效。使用 Clojure

    (defn foo [x] (let [aa @a] (swap! a (fn [&args] x)) aa))
    
    (def a (atom 0))
    (map foo (range 1 20))
    

    相比:

    (def a (atom 0))
    (pmap foo (range 1 20))
    

    【讨论】:

    • @Jyd 我添加了关于原子的简短评论。我想你会发现它很有用。
    • 我非常感谢您为此付出的努力,我不确定在 clojure 开发中使用可变数据结构,正如建议的那样,如果我改变数据结构,它就没有办法了,只是试图避免使用可变数据,尤其是现在我处于早期的 clojure 开发阶段。或者你有什么建议?
    • 快速跟进,我注意到["5" "16" ""] ["07" "" "36"] ["07" "" "34"] 将产生["5" "16" ""] ["07" "16" "36"] ["07" "" "34"] 而不是["5" "16" ""] ["07" "16" "36"] ["07" "16" "34"],除非我们当然会再次重复该过程以填补空白
    • 可以支持任意的look-behind,但是需要一些努力。您应该收集 last-not-empty 进行广播,而不是最后一个元素。然后你必须更正它以处理所有空列的分区(从前一个分区移动最后一个非空值)。最后,您可以使用递归函数代替分区+映射,该函数将使用最后一个非空填充列。
    猜你喜欢
    • 2022-01-05
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-19
    • 2015-11-01
    • 2017-07-15
    • 2021-11-19
    相关资源
    最近更新 更多