【发布时间】:2021-10-07 02:17:47
【问题描述】:
我正在处理 kafka 主题并尝试使用 pyspark 在我的本地计算机中创建一个 readStream。
我已经通过 home-brew 安装了 spark 通过以下命令 brew install apache-spark
我学习了很多教程,但无法到达任何地方。
我还尝试了将 kafka 与 Spark 连接的 Guid -> https://spark.apache.org/docs/2.2.0/structured-streaming-kafka-integration.html。
但这也无济于事。
以下是我将 pyspark 与 Confluent Kafka 主题连接起来的代码
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("sparkkafka").config("spark.master", "local[*]").getOrCreate()
df = spark.readStream.format("kafka")\
.option("kafka.bootstrap.servers", "--xxx--:--xx--") \
.option("subscribe", "NAME OF THE TOPIC") \
.option("startingOffsets", "latest") \
.option("security.protocol", "some protocol") \
.option("mechanisms", "PLAIN") \
.option("[protocol]username", "XXX-username-XXX") \
.option("[protocol]password", "---xxx--password----") \
.option("schema.registry.url", "--- scheme registry url ---") \
.option("basic.auth.credentials.source", "auth source") \
.option("basic.auth.user.info", "info of user") \
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
print(df)
我尝试了两种方式来执行这段代码。
$: python3 fileName$: pyspark --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.4,org.apache.spark:spark-avro_2.11:2.4.0
这两件事都不行。
如果有人已经尝试连接 confluent-kafka 和 pyspark。对于近乎实时的流媒体,您能否指导我一些步骤或一些参考,以便我解决这个问题。
提前致谢
【问题讨论】:
-
homebrew install latest apache-spark 所以它是 3.1.2,所以你应该使用的包是
org.apache.spark:spark-sql-kafka-0-10_2.11:3.1.2。 -
@pltc 感谢您的建议。我帮了我很多。
-
Spark-structured-streaming 只读取字节,顺便说一下,因此您的架构注册表属性被忽略,因此“Confluent”不是这里的问题
标签: python apache-spark pyspark apache-kafka