【发布时间】:2020-10-07 00:37:36
【问题描述】:
我正在 Spark Structured Streaming 中处理一个 kafka JSON 流。作为微批次处理,我可以将累加器与流数据帧一起使用吗?
LongAccumulator longAccum = new LongAccumulator("my accum");
Dataset<Row> df2 = df.filter(output.col("Called number").equalTo("0860"))
.groupBy("Calling number").count();
// put row counter to accumulator for example
df2.javaRDD().foreach(row -> {longAccumulator.add(1);})
抛出
Exception in thread "main" org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start();;
。我也很困惑以这种方式使用累加器。将数据帧转换为 RDD 看起来很奇怪且不必要。我可以在没有 RDD 和 foreach() 的情况下完成吗?
根据例外,我从源数据帧中删除了 foreach 并在 writeStream().foreachBatch() 中完成
StreamingQuery ds = df2
.writeStream().foreachBatch( (rowDataset, aLong) -> {
longAccum.add(1);
log.info("accum : " + longAccum.value());
})
.outputMode("complete")
.format("console").start();
它正在工作,但我在日志中没有值,并且在 GUI 中看不到累加器。
【问题讨论】:
标签: java apache-spark streaming accumulator