【问题标题】:How to specify the group id of kafka consumer for spark structured streaming?如何为火花结构化流指定kafka消费者的组ID?
【发布时间】:2020-11-21 23:39:53
【问题描述】:

我想在同一个 emr 集群中运行 2 个 spark 结构化流作业,以使用同一个 kafka 主题。两个作业都处于运行状态。但是,只有一项工作可以获取 kafka 数据。我对 kafka 部分的配置如下。

        .format("kafka")
        .option("kafka.bootstrap.servers", "xxx")
        .option("subscribe", "sametopic")
        .option("kafka.security.protocol", "SASL_SSL")
          .option("kafka.ssl.truststore.location", "./cacerts")
          .option("kafka.ssl.truststore.password", "changeit")
          .option("kafka.ssl.truststore.type", "JKS")
          .option("kafka.sasl.kerberos.service.name", "kafka")
          .option("kafka.sasl.mechanism", "GSSAPI")
        .load()

我没有设置 group.id。我猜两个工作中的相同组 id 被用来导致这个问题。但是,当我设置 group.id 时,它抱怨“用户指定的消费者组不用于跟踪偏移量。”。解决这个问题的正确方法是什么?谢谢!

【问题讨论】:

标签: apache-spark apache-spark-sql spark-streaming spark-streaming-kafka


【解决方案1】:

您需要运行 Spark v3。

来自https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html

kafka.group.id

从 Kafka 读取时在 Kafka 消费者中使用的 Kafka 组 ID。 请谨慎使用。默认情况下,每个查询都会生成一个唯一的组 用于读取数据的 ID。这确保了每个 Kafka 源都有自己的 不受其他任何干扰的消费群体 消费者,因此可以读取其所有分区 订阅的主题。在某些场景下(例如,基于 Kafka 组的 授权),您可能希望使用特定的授权组 ID 读取数据。您可以选择设置组 ID。但是,这样做 非常小心,因为它可能导致意外行为。同时 运行查询(批处理和流式处理)或具有相同的源 组 ID 可能会相互干扰,导致每个查询 只读取部分数据。这也可能在查询时发生 快速连续启动/重新启动。为了尽量减少此类问题,请设置 Kafka 消费者会话超时(通过设置选项 “kafka.session.timeout.ms”)非常小。设置时,选项 “groupIdPrefix”将被忽略。

【讨论】:

  • 谢谢,我会试试 spark v3。 Spark v3 现在是否已经集成到 EMR 中?
猜你喜欢
  • 2018-07-18
  • 1970-01-01
  • 2018-11-06
  • 2019-09-20
  • 2020-02-14
  • 2019-07-22
  • 2019-08-16
相关资源
最近更新 更多