【问题标题】:How do I reduce a spark dataframe from kafka and collect the result?如何减少来自 kafka 的 spark 数据帧并收集结果?
【发布时间】: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


    【解决方案1】:

    您只能在流式数据帧上使用writeStream。我不确定您是否打算拥有此流数据帧。如果您删除readStream 并改用read,您可能会解决此问题!

    【讨论】:

    • 就是这样。谢谢!
    猜你喜欢
    • 2015-11-07
    • 2017-08-25
    • 2019-04-28
    • 1970-01-01
    • 1970-01-01
    • 2017-08-07
    • 1970-01-01
    • 2022-09-22
    • 1970-01-01
    相关资源
    最近更新 更多