【发布时间】:2018-01-16 14:13:14
【问题描述】:
尽管我使用的是withWatermark(),但我在运行 Spark 作业时收到以下错误消息:
线程“main”中的异常 org.apache.spark.sql.AnalysisException:当流式 DataFrames/DataSets 上存在流式聚合时,不支持附加输出模式而没有水印;;
从我在programming guide 中看到的内容来看,这完全符合预期用途(以及示例代码)。有谁知道可能出了什么问题?
提前致谢!
相关代码(Java 8、Spark 2.2.0):
StructType logSchema = new StructType()
.add("timestamp", TimestampType)
.add("key", IntegerType)
.add("val", IntegerType);
Dataset<Row> kafka = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", brokers)
.option("subscribe", topics)
.load();
Dataset<Row> parsed = kafka
.select(from_json(col("value").cast("string"), logSchema).alias("parsed_value"))
.select("parsed_value.*");
Dataset<Row> tenSecondCounts = parsed
.withWatermark("timestamp", "10 minutes")
.groupBy(
parsed.col("key"),
window(parsed.col("timestamp"), "1 day"))
.count();
StreamingQuery query = tenSecondCounts
.writeStream()
.trigger(Trigger.ProcessingTime("10 seconds"))
.outputMode("append")
.format("console")
.option("truncate", false)
.start();
【问题讨论】:
标签: java apache-spark spark-structured-streaming