【发布时间】:2020-04-11 15:34:39
【问题描述】:
我正在 PySpark 中编写一个 Spark 结构化流应用程序,以从 Confluent Cloud 中的 Kafka 读取数据。 spark readstream() 函数的文档太浅了,并且没有在可选参数部分特别是在身份验证机制部分指定太多。我不确定哪个参数出错并导致连接崩溃。任何有 Spark 经验的人都可以帮助我开始这种连接吗?
必填参数
> Consumer({'bootstrap.servers': > 'cluster.gcp.confluent.cloud:9092', > 'sasl.username':'xxx', > 'sasl.password': 'xxx', > 'sasl.mechanisms': 'PLAIN', > 'security.protocol': 'SASL_SSL', > 'group.id': 'python_example_group_1', > 'auto.offset.reset': 'earliest' })
这是我的 pyspark 代码:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "cluster.gcp.confluent.cloud:9092") \
.option("subscribe", "test-topic") \
.option("kafka.sasl.mechanisms", "PLAIN")\
.option("kafka.security.protocol", "SASL_SSL")\
.option("kafka.sasl.username","xxx")\
.option("kafka.sasl.password", "xxx")\
.option("startingOffsets", "latest")\
.option("kafka.group.id", "python_example_group_1")\
.load()
display(df)
但是,我不断收到错误消息:
kafkashaded.org.apache.kafka.common.KafkaException: 失败 构建kafka消费者
DataBrick Notebook- 用于测试
文档
【问题讨论】:
-
您是否尝试过使用其他消费者? spark.apache.org/docs/2.3.1/…
-
@cricket_007,还没有,因为要求之一是流-流连接,它只支持结构化流。 Direct Stream 在加入两个流时有限制。因此,这就是我需要使用
readStream()的原因。在线资源只关注 DStream。很难找到答案
标签: apache-spark pyspark apache-kafka spark-streaming