【发布时间】:2019-11-29 01:21:19
【问题描述】:
我有一个从 kafka 流创建的数据框。我想将其减少为单个值,然后在我的程序中使用该单个值。
```scala
import sparkSession.implicits._
val df = sparkSession
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", ...)
.option("subscribe", "theTopic")
.load()
val result = df
.selectExpr("CAST(value AS STRING) as json")
.map(json => getAnInt(json))
.reduce { (x, y) =>
if (x > y) x else y
}
someOtherFunction(result)
```
我希望将流减少到一个值,然后我可以在我的程序的其余部分中使用它。相反,它失败了:
org.apache.spark.sql.AnalysisException: 带有流源的查询必须使用 writeStream.start();; 卡夫卡 在 org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.throwError(UnsupportedOperationChecker.scala:389) 在 org.apache.spark.sql.catalyst.analysis.U...
【问题讨论】:
标签: scala apache-spark apache-kafka reduce