【发布时间】:2021-03-28 18:47:32
【问题描述】:
我正在使用带有 Kafka 的 Spark 结构化流,并且主题已被订阅为模式:
option("subscribePattern", "topic.*")
// Subscribe to a pattern
val df = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribePattern", "topic.*")
.load()
一旦我开始工作并列出了一个新主题,比如topic.new_topic,该工作就不会自动开始收听新主题,它需要重新启动。
有没有办法在不重新启动作业的情况下自动订阅新模式?
火花:3.0.0
【问题讨论】:
-
自 3.0.0 起,消费者缓存可用,这可能允许我们从新主题中读取数据 - spark.apache.org/docs/latest/…
标签: apache-spark apache-kafka spark-structured-streaming spark-kafka-integration