【发布时间】:2020-04-03 20:01:56
【问题描述】:
我正在尝试通过 Spark 结构化流从 Kafka 读取数据。但是,在 Spark 2.4.0 中,您无法为流设置组 ID(请参阅 How to set group.id for consumer group in kafka data source in Structured Streaming?)。
但是,由于没有设置,spark 只是生成组 ID,我被困在 GroupAuthorizationException:
19/12/10 15:15:00 ERROR streaming.MicroBatchExecution: Query [id = 747090ff-120f-4a6d-b20e-634eb77ac7b8, runId = 63aa4cce-ad72-47f2-80f6-e87947b69685] terminated with error
org.apache.kafka.common.errors.GroupAuthorizationException: Not authorized to access group: spark-kafka-source-d2420426-13d5-4bda-ad21-7d8e43ebf518-1874352823-driver-2
任何想法如何绕过这个请?有趣的是,我可以通过 kafka-console-consumer.sh 读取这些数据,我可以在 .properties 文件中传递组 ID。
抛出异常的代码:
val df = spark
.readStream
.format("kafka")
.option("subscribe", "topic")
.option("startingOffsets", "earliest")
.option("kafka.group.id", "idThatShouldBeUsed")
.option("kafka.bootstrap.servers", "server")
.option("kafka.security.protocol", "SASL_SSL")
.option("kafka.sasl.mechanism", "PLAIN")
.option("kafka.ssl.truststore.location", "/location)
.option("kafka.ssl.truststore.password", "pass")
.option("kafka.sasl.jaas.config", """jaasToUse""")
.load()
.writeStream
.outputMode("append")
.format("console")
.option("startingOffsets", "earliest")
.start().awaitTermination()
【问题讨论】:
-
组 ID 不应决定身份验证。 JKS 文件和 JAAS 属性应该
-
好吧,似乎 - 在授予组权限时使用通配符可以解决相同的问题 (stackoverflow.com/questions/48545215/…)。但是,我不能更改这些 Kafka 设置。
-
我的印象是那些是“用户组”,而不是消费者的“组 id”。顺便说一下,Authorizer 是可插拔的,但您必须与 Kafka 管理员一起调整这些设置
标签: scala apache-spark apache-kafka spark-streaming spark-structured-streaming