【问题标题】:How to set group.id for consumer group in kafka data source in Structured Streaming?如何在结构化流的kafka数据源中为消费者组设置group.id?
【发布时间】:2019-08-16 17:31:12
【问题描述】:

我想使用 Spark Structured Streaming 从安全的 kafka 中读取数据。这意味着我需要强制使用特定的 group.id。但是,正如文档中所述,这是不可能的。 尽管如此,在databricks 文档https://docs.azuredatabricks.net/spark/latest/structured-streaming/kafka.html#using-ssl 中,它说这是可能的。这只是指天蓝色集群​​吗?

此外,通过查看 apache/spark repo https://github.com/apache/spark/blob/master/docs/structured-streaming-kafka-integration.md 的 master 分支的文档,我们可以了解到此类功能旨在在以后的 spark 版本中添加。你知道这样一个稳定版本的任何计划,这将允许设置消费者 group.id 吗?

如果没有,Spark 2.4.0 是否有任何变通方法可以设置特定的使用者 group.id?

【问题讨论】:

    标签: apache-spark apache-kafka spark-structured-streaming spark-kafka-integration


    【解决方案1】:

    目前 (v2.4.0) 这是不可能的。

    您可以在 Apache Spark 项目中检查以下几行:

    https://github.com/apache/spark/blob/v2.4.0/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala#L81 - 生成 group.id

    https://github.com/apache/spark/blob/v2.4.0/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala#L534 - 在属性中设置它,用于创建KafkaConsumer

    在 master 分支中你可以找到修改,可以设置 prefix 或特定的 group.id

    https://github.com/apache/spark/blob/master/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala#L83 - 根据组前缀 (groupidprefix) 生成 group.id

    https://github.com/apache/spark/blob/master/external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala#L543 - 设置之前生成的 groupId,如果 kafka.group.id 没有在属性中传递

    【讨论】:

    • 感谢您的回复,您知道我将如何实施这些修改后的类吗?从包中构建一个 jar 并使用 spark-submit 添加该 jar 就足够了吗?
    • @PanagiotisFytas,您可以在 apache spark 的 master 分支中检查代码。我认为删除以下行(github.com/apache/spark/blob/v2.4.0/external/kafka-0-10-sql/src/…)并构建并添加 jar 到 spark-submit 并通过 option 传递 kafka.group.id 属性就足够了@
    • 我通过从 master 分支向 2.4.0 分支添加一些提交来完成这项工作。你可以查看我的 fork github.com/PanagiotisFytas/spark。您可以通过调用 [external/kafka-0-10-sql/mvn -DskipTests package] 使用自定义连接器构建胖 jar。我不确定这种方法有多安全,因为我还没有完全测试过。
    • @PanagiotisFytas,我会非常小心。肯定不推荐在生产中使用这种方法。然而它只是对 kafka 流部分的修改。
    • 到目前为止,Dstream 是生产实施的方式。
    【解决方案2】:

    自 Spark 3.0.0

    根据Structured Kafka Integration Guide,您可以提供ConsumerGroup 作为选项kafka.group.id

    val df = spark
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
      .option("subscribe", "topic1")
      .option("kafka.group.id", "myConsumerGroup")
      .load()
    

    但是,Spark 不会提交任何偏移量,因此您的 ConsumerGroups 的偏移量不会存储在 Kafka 的内部主题 __consumer_offsets 中,而是存储在 Spark 的检查点文件中。

    能够设置group.id 意味着处理Kafka 的最新功能Authorization using Role-Based Access Control,您的ConsumerGroup 通常需要遵循命名约定。

    讨论并解决了一个 Spark 3.x 应用程序设置kafka.group.id 的完整示例here

    【讨论】:

      【解决方案3】:

      【讨论】:

      • 它有一条警告信息,我们不应该使用 group.id。如果您将日志打开为 WARN 级别,您会看到这一点。
      【解决方案4】:

      Structured Streaming guide 似乎对此非常明确:

      请注意,以下 Kafka 参数 无法设置,而 Kafka source 或 sink 会抛出异常:

      group.id:Kafka 源将为每个查询创建一个唯一的组 id 自动地。

      auto.offset.reset:设置源选项 startingOffsets 指定从哪里开始。

      【讨论】:

      • 我知道。当时,我要求解决方法,这是可能的。我设法使它与建议的解决方案一起工作。
      • @PanagiotisFytas 您能否显示代码(作为单独的答案)以帮助我们更好地理解它?那会非常有帮助。谢谢。
      猜你喜欢
      • 2019-08-09
      • 2020-11-21
      • 2017-10-19
      • 2019-07-22
      • 2019-10-05
      • 2019-09-24
      • 2018-11-23
      相关资源
      最近更新 更多