【发布时间】:2019-12-10 22:18:48
【问题描述】:
我正在使用 spark 结构化流从 kafka 读取数据。
val readStreamDF = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", config.getString("kafka.source.brokerList"))
.option("startingOffsets", config.getString("kafka.source.startingOffsets"))
.option("subscribe", config.getString("kafka.source.topic"))
.load()
基于从 kafka 读取的消息中的uid,我必须对外部源进行 api 调用并获取数据并写回另一个 kafka 主题。
为此,我使用自定义 foreach 编写器并处理每条消息。
import spark.implicits._
val eventData = readStreamDF
.select(from_json(col("value").cast("string"), event).alias("message"), col("timestamp"))
.withColumn("uid", col("message.eventPayload.uid"))
.drop("message")
val q = eventData
.writeStream
.format("console")
.foreach(new CustomForEachWriter())
.start()
CustomForEachWriter 进行 API 调用并根据给定的uid 从服务中获取结果。结果是一个 id 数组。然后这些 id 再次通过 kafka 生产者写回另一个 kafka 主题。
有 30 个 kafka 分区,我使用以下配置启动了 spark
num-executors = 30
executors-cores = 3
executor-memory = 10GB
但火花作业仍然开始滞后,无法跟上传入的数据速率。
传入数据速率约为每秒 10K 条消息。处理单个消息的平均时间为 100 毫秒。
我想了解在结构化流的情况下 spark 是如何处理这个的。 在结构化流的情况下,有一个专用的执行器负责从 kafka 的所有分区读取数据。 该执行者是否根据否分配任务。 kafka 中的分区。 批处理中的数据按顺序处理。如何使其并行处理以最大化吞吐量。
【问题讨论】:
标签: scala apache-spark spark-structured-streaming