【发布时间】:2021-01-14 23:12:53
【问题描述】:
我正在开发一个实时流应用程序,该应用程序从 Kafka 代理轮询数据,并且我正在调整以前默认使用 Spark 结构化流处理的代码(使用微批处理)。但是,我不知道如何使用连续流而不是微批处理流来获得类似的行为。这是一段有效的代码:
query = df.writeStream \
.foreachBatch(foreach_batch_func) \
.start()
这是我迄今为止尝试的连续流式传输:
query = df \
.writeStream \
.foreach(example_func) \
.trigger(continuous = '1 second') \
.start()
应用弹出如下错误:
连续执行不支持
org.apache.spark.sql.execution.streaming.continuous.ContinuousDataSourceRDD.compute(ContinuousDataSourceRDD.scala:76)处的任务重试
我正在使用带有 Scala 2.12、Kafka 2.6.0 的 Spark (pyspark) 3.0.1
当我提交应用程序时,我正在添加 jar org.apache.spark:spark-sql-kafka-0-10_2.12:3.0.1。
【问题讨论】:
标签: apache-spark pyspark apache-kafka spark-structured-streaming