【发布时间】: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