【问题标题】:Error when writing Spark Structured Streaming example in Clojure在 Clojure 中编写 Spark 结构化流示例时出错
【发布时间】:2017-10-10 01:07:34
【问题描述】:

我正在尝试在 Clojure 中重写 Spark Structured Streaming 示例。

示例用 Scala 编写如下:

https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html

(ns flambo-example.streaming-example
  (:import [org.apache.spark.sql Encoders SparkSession Dataset Row]
           [org.apache.spark.sql.functions]
           ))

(def spark
  (->
   (SparkSession/builder)
   (.appName "sample")
   (.master "local[*]")
   .getOrCreate)  
  )


(def lines
  (-> spark
      .readStream
      (.format "socket")
      (.option "host" "localhost")
      (.option "port" 9999)
      .load      
      )
  )

(def words
  (-> lines
      (.as (Encoders/STRING))      
      (.flatMap #(clojure.string/split  % #" " ))      
      ))

以上代码导致以下异常。

;;由 java.lang.IllegalArgumentException 引起 ;;未找到匹配方法:类的 flatMap ;; org.apache.spark.sql.Dataset

如何避免错误?

【问题讨论】:

    标签: scala apache-spark clojure spark-structured-streaming flambo


    【解决方案1】:

    您必须遵循签名。 JavaDataset API 提供了两种Dataset.flatMap 的实现,一种采用scala.Function1

    def flatMap[U](func: (T) ⇒ TraversableOnce[U])(implicit arg0: Encoder[U]): Dataset[U] 
    

    第二个使用 Spark 自己的o.a.s.api.java.function.FlatMapFunction

    def flatMap[U](f: FlatMapFunction[T, U], encoder: Encoder[U]): Dataset[U] 
    

    前一个对你来说没什么用,但你应该可以使用后一个。对于 RDD API flambo uses macros to create Spark friendly adapters,可以使用 flambo.api/fn 访问 - 我不确定这些是否可以直接使用 Datasets,但如果需要,您应该能够调整它们。

    由于您不能依赖隐式Encoders,您还必须提供与返回类型匹配的显式编码器。

    总的来说,你需要一些东西:

    (def words
      (-> lines
        (.as (Encoders/STRING))      
        (.flatMap f e)      
      ))
    

    其中f 实现FlatMapFunction 并且eEncoder。一个示例实现:

    (def words
      (-> lines
          (.as (Encoders/STRING))      
          (.flatMap
            (proxy [FlatMapFunction] [] 
              (call [s] (.iterator (clojure.string/split s #" ")))) 
            (Encoders/STRING))))
    

    但我想有可能找到更好的。

    实际上,我会避免输入Dataset,而是专注于DataFrame (Dataset[Row])。

    【讨论】:

      猜你喜欢
      • 2020-01-30
      • 2018-03-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-05-04
      • 2019-04-06
      • 2018-10-06
      • 1970-01-01
      相关资源
      最近更新 更多