【问题标题】:How many Kafka consumers does a streaming query use for execution?流式查询使用多少个 Kafka 消费者来执行?
【发布时间】:2019-05-05 09:56:59
【问题描述】:

令我惊讶的是,Spark 仅使用一个 Kafka 消费者来使用来自 Kafka 的数据,并且该消费者在驱动程序容器中运行。我更希望看到,Spark 创建的消费者数量与主题中的分区数量一样多,并在执行器容器中运行这些消费者。

例如,我有一个主题 events 有 5 个分区。我启动了我的 Spark Structured Streaming 应用程序,该应用程序使用该主题并写入 HDFS 上的 Parquet。该应用程序有 5 个执行者。 在检查 Spark 创建的 Kafka 消费者组时,我发现只有一个消费者负责所有 5 个分区。这个消费者正在使用驱动程序的机器上运行:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group spark-kafka-source-08e10acf-7234-425c-a78b-3552694f22ef--1589131535-driver-0

TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CONSUMER-ID                                     HOST            CLIENT-ID
events          2          -               0               -               consumer-1-8c3d806d-eb1e-4536-97d5-7c9d19582942 /192.168.100.147  consumer-1
events          1          -               0               -               consumer-1-8c3d806d-eb1e-4536-97d5-7c9d19582942 /192.168.100.147  consumer-1
events          0          -               0               -               consumer-1-8c3d806d-eb1e-4536-97d5-7c9d19582942 /192.168.100.147  consumer-1
events          4          -               0               -               consumer-1-8c3d806d-eb1e-4536-97d5-7c9d19582942 /192.168.100.147  consumer-1
events          3          -               0               -               consumer-1-8c3d806d-eb1e-4536-97d5-7c9d19582942 /192.168.100.147  consumer-1

在检查了所有 5 个 executor 的日志后,我发现只有一个在忙于将消费数据写入 HDFS 上的 Parquet 位置。其他 4 个闲置。

这很奇怪。我的期望是 5 个执行器应该并行使用来自 5 个 Kafka 分区的数据并在 HDFS 上并行写入。这是否意味着驱动程序使用来自 Kafka 的数据并将其分发给执行程序?它看起来像一个瓶颈。

UPDATE 1 我尝试将 repartition(5) 添加到流数据帧:

spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "brokerhost:9092")
    .option("subscribe", "events")
    .option("startingOffsets", "earliest")
    .load()
    .repartition(5)

之后,我看到所有 5 个执行程序都将数据写入 HDFS(根据他们的日志)。尽管如此,我在 Kafka 主题的所有 5 个分区上只看到一个消费者(驱动程序)。

更新 2 Spark 版本 2.4.0。以下是提交申请的命令:

spark-submit \
--name "Streaming Spark App" \
--master yarn \
--deploy-mode cluster \
--conf spark.yarn.maxAppAttempts=1 \
--conf spark.executor.instances=5 \
--conf spark.sql.shuffle.partitions=5 \
--class example.ConsumerMain \
"$jar_file"

【问题讨论】:

  • 不确定是否可以在火花流中拥有多个 Kafka 主题消费者。但是,spark 的默认分区计数而 KafkaUtils.createDirectStream 等于 Kafka 主题的分区计数。因此,在您的情况下,所有 5 个执行程序都应该将数据写入 HDFS,而无需重新分区,从而降低重新洗牌成本。所以建议使用 KafkaUtils.createDirectStream 而不是 spark.readStream。
  • @aavos “而 KafkaUtils.createDirectStream 等于 Kafka 主题的分区数。” 是关于 Spark Streaming,但 OP 询问的是 Spark Structured Streaming。它们是不同的流媒体引擎。
  • 什么是 Spark 版本?你如何spark-submit 应用程序?
  • 用 Spark 版本和提交命令更新了帖子
  • 对此有任何反馈吗?我在使用 Spark 2.4.4 时遇到了同样的问题

标签: apache-kafka spark-structured-streaming


【解决方案1】:

根据结构化流的文档,我可以看到它被提及为在执行程序上创建的消费者,Consumer Caching

【讨论】:

    猜你喜欢
    • 2019-11-18
    • 2013-08-30
    • 2019-04-07
    • 2018-08-11
    • 2016-04-13
    • 1970-01-01
    • 1970-01-01
    • 2018-09-29
    • 1970-01-01
    相关资源
    最近更新 更多