【问题标题】:How to read json string from kafka topic into pyspark dataframe?如何将 json 字符串从 kafka 主题读取到 pyspark 数据帧中?
【发布时间】: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


【解决方案1】:

您缺少 main documentation 中提到的 kafka 包

./bin/spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 ...

确保此处列出的 3.1.2 与您自己的 Spark 版本匹配

【讨论】:

  • 非常感谢您的回答,它帮助我解决了问题,但我遇到了新问题。我更新了问题;请再检查一次
猜你喜欢
  • 2021-03-01
  • 2020-04-03
  • 2020-12-09
  • 2018-05-07
  • 2018-01-23
  • 2021-07-03
  • 2023-03-30
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多