【发布时间】:2018-12-26 07:02:08
【问题描述】:
我有这个代码
val counter = event_stream
.withWatermark("timestamp", "5 minutes")
.groupBy(
window($"timestamp", "10 minutes", "5 minutes"),
$"value")
.agg(count("value") as "kafka.count",collect_set("topic") as "kafka.topic")
.drop("window")
.withColumnRenamed("value","join_id")
counter.printSchema
val counter1 = event_stream
.groupBy("value")
.count()
// .agg(count("value") as "kafka.count",collect_set("topic") as "kafka.topic")
.withColumnRenamed("value","join_id")
counter1.printSchema()
val result_stream = event_stream.join(counter,$"value" === $"join_id")
.drop("key")
.drop("value")
.drop("partition")
.drop("timestamp")
.drop("join_id")
.drop("timestampType")
.drop("offset")
// .withColumnRenamed("count(value)", "kafka.count")
.withColumnRenamed("topic","kafka.topic")
result_stream.printSchema()
KafkaSink.write(counter, topic_produce)
// KafkaSink.writeToConsole(result_stream, topic_produce)
如果我将它发送到我使用 Outputmode.Complete 的控制台,它可以正常工作,但是当我使用 OutputMode.Append 时。在上面发送不同的流式查询时,它会给出不同的错误。
这是我的写函数
private def writeStream(df:DataFrame, topic:String): StreamingQuery = {
df
.writeStream
.format("kafka")
.option("kafka.bootstrap.servers", KafkaUtils.kafkaServers)
.option("topic", topic)
.option("checkpointLocation", KafkaUtils.checkPointDir)
.outputMode(OutputMode.Append())
.start()
}
我收到此错误
java.lang.IllegalArgumentException: Expected e.g. {"topicA":{"0":23,"1":-1},"topicB":{"0":-2}}, got 1
{"path":"file:///home/ukaleem/Documents/freenet/Proto2/src/main/resource/events-identification-carrier-a.txt","timestamp":1530198790000,"batchId":0}
为什么会出现这个错误?
第 2 部分:如果我从上面的代码开始
val result_stream = event_stream.join(counter,$"value" === $"join_id")
KafkaSink.write(result_stream, topic_produce)
我收到此错误
java.lang.AssertionError: assertion failed
at scala.Predef$.assert(Predef.scala:156)
at org.apache.spark.sql.execution.streaming.OffsetSeq.toStreamProgress(OffsetSeq.scala:42)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.org$apache$spark$sql$execution$streaming$MicroBatchExecution$$populateStartOffsets(MicroBatchExecution.scala:185)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply$mcV$sp(MicroBatchExecution.scala:124)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply(MicroBatchExecution.scala:121)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply(MicroBatchExecution.scala:121)
at org.apache.spark.sql.execution.streaming.ProgressReporter$class.reportTimeTaken(ProgressReporter.scala:271)
at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:58)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1.apply$mcZ$sp(MicroBatchExecution.scala:121)
at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:56)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:117)
at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:279)
at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:189)
Exception in thread "main" org.apache.spark.sql.streaming.StreamingQueryException: assertion failed
这两种情况都对我有用。但我在这两个方面都遇到了错误。
编辑:我解决了第一部分。但是,仍然需要第二个。
【问题讨论】:
-
你能回答你自己的问题来描述“编辑:我解决了第一部分。但仍然需要第二部分。”即使第二个问题还没有解决?我什至建议提出两个单独的问题,并将第一部分作为单独的问题回答。 WDYT?
-
什么时候得到“java.lang.AssertionError: assertion failed”?启动应用程序时是否总是发生这种情况?您是否使用任何
checkpointLocation选项?您是否更改了中间的结构化查询以使用不同的来源? -
是的,我都能解决这两个问题。您面临什么问题?
-
@Sam 我解决了这个问题,请考虑answering your own question。
-
@JacekLaskowski 我遇到了同样的问题。仅当我使用现有的 checkpointLocation 时才会发生这种情况。我更改了代码并想在它停止的地方启动应用程序,但 spark 会引发此错误。你有什么解决办法吗?
标签: apache-spark apache-kafka spark-structured-streaming