【发布时间】:2020-10-02 15:12:55
【问题描述】:
我正在尝试编写一个 Spark Structured Streaming 作业,该作业从多个 Kafka 主题(可能是 100 个)中读取,并根据主题名称将结果写入 S3 上的不同位置。我已经开发了这个 sn-p 的代码,它当前从多个主题中读取并将结果输出到控制台(基于循环)并且它按预期工作。但是,我想了解性能影响是什么。这是推荐的方法吗?不建议有多个 readStream 和 writeStream 操作吗?如果是,推荐的方法是什么?
my_topics = ["topic_1", "topic_2"]
for i in my_topics:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", bootstrap_servers) \
.option("subscribePattern", i) \
.load() \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
output_df = df \
.writeStream \
.format("console") \
.option("truncate", False) \
.outputMode("update") \
.option("checkpointLocation", "s3://<MY_BUCKET>/{}".format(i)) \
.start()
【问题讨论】:
-
为什么要为每个主题设置不同的 checkpointLocation 位置,所有主题都可以使用一个?
-
Kafka Connect 通常是 Kafka -> S3 的更好方法。如果有用的话,我可以提供一个答案。
-
@Srinivas 如果我需要通过清除检查点位置来重新启动/重置特定主题,最好有单独的检查点位置以避免与其他主题的检查点?
-
@RobinMoffatt 我已经探索了使用 Kafka Connect 的选项,但是,我想使用 Spark Structured Streaming 来扩展接收器的数量。
-
(Kafka Connect 可以处理正则表达式主题列表,如果您担心的话。)
标签: apache-spark pyspark apache-kafka spark-structured-streaming