【发布时间】:2020-09-28 03:16:00
【问题描述】:
我想构建一个从多个 Kafka 主题(数量随时间变化)读取的 Spark 流式传输管道。我打算使用Spark Structured Streaming + Kafka Integration Guide 中列出的两个选项之一来停止流式传输作业,添加/删除新主题,并在需要更新流式传输作业中的主题时再次启动该作业:
# Subscribe to multiple topics
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribe", "topic1,topic2") \
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
# Subscribe to a pattern
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribePattern", "topic.*") \
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
经过进一步调查,我注意到Spark Structured Streaming Programming Guide 中的以下几点,并试图了解为什么“不允许”更改输入源的数量:
输入源的数量或类型(即不同的源)的变化:这是不允许的。
“不允许”的定义(也来自Spark Structured Streaming Programming Guide):
术语不允许意味着您不应该进行指定的更改,因为重新启动的查询很可能会因不可预知的错误而失败。 sdf 表示使用 sparkSession.readStream 生成的流式 DataFrame/Dataset。
我的理解是 Spark Structured Streaming 实现了自己的checkpointing mechanism:
如果出现故障或故意关闭,您可以恢复之前查询的进度和状态,并从中断处继续。这是使用检查点和预写日志完成的。您可以使用检查点位置配置查询,并且查询会将所有进度信息(即每个触发器中处理的偏移范围)和正在运行的聚合(例如快速示例中的字数)保存到检查点位置。此检查点位置必须是 HDFS 兼容文件系统中的路径,并且可以在启动查询时在 DataStreamWriter 中设置为选项。
有人可以解释为什么“不允许”更改来源的数量吗?我认为这将是检查点机制的好处之一。
【问题讨论】:
标签: apache-spark pyspark apache-kafka spark-structured-streaming