【发布时间】:2021-08-21 22:02:27
【问题描述】:
我正在尝试将来自 Kafka 主题的 json 消息读入 PySpark 数据帧。我的第一反应是这样的:
consumer = KafkaConsumer(TOPIC_NAME,
consumer_timeout_ms=9000,
bootstrap_servers=BOOTSTRAP_SERVER,
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id=str(uuid4()),
value_deserializer=lambda x: x.decode("utf-8"))
message_lst = []
for message in consumer:
message_str = message.value.replace('\\"', "'").replace("\n", "").replace("\r", "")
message_dict = json.loads(message_str)
message_lst.append(message_dict)
messages_json = sc.parallelize(message_lst)
messages_df = sqlContext.read.json(messages_json)
我想知道有没有办法使用 Spark 结构化流或类似的东西来获取相同的数据帧。有人可以帮忙吗? UPD:我对结构化流的尝试是这样的:
df = spark \
.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", f"{BOOTSTRAP_SERVER}") \
.option("subscribe", TOPIC_NAME) \
.load()
它退出并出现以下错误:
pyspark.sql.utils.AnalysisException: Failed to find data source: Kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".
UPD:我阅读了异常文本中所述的指南,它建议安装此库“spark-sql-kafka-0-10_2.12”,但我找不到。有人知道吗?
UPD 2:我设法添加了所需的包并尝试了从 kafka 读取消息:
df = spark \
... .readStream \
... .format("kafka") \
... .option("kafka.bootstrap.servers", f"{BOOTSTRAP_SERVER}") \
... .option("subscribe", TOPIC_NAME) \
... .load()
df.writeStream.outputMode("append").format("console").start().awaitTermination()
我使用与以前相同的消费者。这里的问题是它只读取在 start() 调用之后写入的消息。如何读取在给定时间写入的所有消息并将结果作为数据框获取?另外,任何人都可以举一个 load_json() 模式的例子吗?如果我的问题很愚蠢,我很抱歉,但我在 Python 中找不到任何示例。
【问题讨论】:
-
这应该会有所帮助:spark.apache.org/docs/latest/…。它展示了如何使用 Spark Structured Streaming 读取 Kafka 流。在示例中,值被读取为字符串,但您可以使用内置函数
from_json轻松地将它们解释为 json -
所以,您知道结构化流式传输,但不清楚您尝试过什么
-
@OneCricketeer 更新了问题;请再检查一次。
-
找不到是什么意思?这是一个 Maven 包,而不是 Python mvnrepository.com/artifact/org.apache.spark/…
标签: python apache-spark pyspark apache-kafka