【问题标题】:Pyspark not printing any data from kafka stream, not failing either [duplicate]Pyspark没有从kafka流中打印任何数据,也没有失败[重复]
【发布时间】:2021-08-11 13:13:39
【问题描述】:

我是 spark 和 kafka 的新手。使用从免费的 Kafka 服务器提供商 (Cloudkarafka) 创建的 kafka 服务器来使用数据。在运行 pyspark 代码(在 databricks 上)以使用流数据时,流只是保持初始化,并且不获取任何内容。它既不会失败,也不会停止执行,只是一直显示状态为“Stream Initializing”。

代码:

from pyspark.sql.functions import col

kafkaServer="<server>"

editsDF=(spark.readStream
        .format("kafka")
        .option("kafka.bootstrap.servers",kafkaServer)
        .option("sasl.username","<username>")
        .option("sasl.password","<password>")
        .option("group.id", "%s-consumer" % "<username>")
        .option("session.timeout.ms", 6000)
        .option("default.topic.config", {"auto.offset.reset": "smallest"})
        .option('security.protocol', 'SASL_SSL')
        .option('sasl.mechanisms', 'SCRAM-SHA-256')
        .option("subscribe","<topic>")
        .option("startingOffsets","latest")
        .option("maxOffsetsPerTrigger",1000)
        .load()
        .select(col("value").cast("STRING"))
        )


query = editsDF \
    .writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

The status in databricks while running the code:

如果我遗漏了什么,请告诉我。提前致谢。

注意:我已经确保 kafka 服务器能够产生消息并且能够在 python 程序中使用它。但不能在 pyspark 中工作。此外,数据量非常小,因此不会出现性能问题。

编辑:这个建议的函数 display() 仍然没有为这个有问题的 Kafka 服务器打印任何数据,但是当我尝试完全使用另一个 Kafka 服务器时它工作正常。我认为这是因为这个 kafka 服务器(有问题)正在使用 SASL-SCRAM 身份验证,所以可能需要对它进行一些不同的配置。如果您从 Pyspark 连接 SASL Kafka,请提供任何详细信息/链接/示例。谢谢!

【问题讨论】:

  • 你开始后有没有往Kafka topic里放数据?
  • 不应添加 query.awaitTermination() 以不终止您的主进程
  • @AlexOtt 是的,每当我运行代码时,我都会手动将数据发送到 Kafka 主题中。
  • @1pluszara 我添加了 query.awaitTermination(),但代码一直在运行。仍然没有打印任何东西。

标签: apache-spark pyspark apache-kafka databricks spark-structured-streaming


【解决方案1】:

当您使用console 接收器时,它会将数据打印到标准输出(请参阅Spark docs),因此您需要在集群 UI 中检查驱动程序日志以获取生成的数据。

要查看 Databricks 笔记本本身中的数据,您需要使用支持显示结构化流中数据的 display 函数(请参阅 Databricks docs)。所以不是

query = editsDF \
    .writeStream \
    .outputMode("append") \
    .format("console") \
    .start()

你只需要写:

display(editsDF)

您还可以向此函数传递其他选项,例如 checkpointLocationtrigger 等 - 请查看我上面链接的文档。

【讨论】:

  • 谢谢亚历克斯。早些时候我不知道我们可以使用 display() 来查看流数据。但是,问题似乎有所不同。这个 display() 仍然没有为这个有问题的 Kafka 服务器打印任何数据,但是当我尝试完全使用另一个 Kafka 服务器时它工作正常。我认为这是因为 kafka 服务器(有问题)正在使用 SASL 身份验证,所以可能需要对它进行一些不同的配置。如果您有从 Pyspark 连接 SASL Kafka 的信息,请提供任何详细信息/链接。谢谢!
  • 如果您遇到了身份验证问题,那么操作将失败 - 如果它有效,那么它是其他的东西,例如,您发送到错误的主题或其他原因
猜你喜欢
  • 2021-07-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-05-02
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-06-14
相关资源
最近更新 更多