【问题标题】:Spark Structured Streaming with foreach使用 foreach 进行 Spark 结构化流式处理
【发布时间】: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


    【解决方案1】:

    我认为CustomForEachWriter writer 将处理数据集的单行/记录。如果您使用的是2.4 版本的 Spark,您可以试验一下foreachBatch。但它正在进化中。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-08-11
      • 1970-01-01
      • 2019-10-03
      • 2019-08-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多