【问题标题】:Read data from Kafka and print to console with Spark Structured Sreaming in Python使用 Python 中的 Spark Structured Streaming 从 Kafka 读取数据并打印到控制台
【发布时间】:2021-04-10 22:16:41
【问题描述】:

我在 Ubuntu 20.04 中有 kafka_2.13-2.7.0。我运行 kafka 服务器和 zookeeper,然后创建一个主题并通过nc -lk 9999 在其中发送一个文本文件。该主题充满了数据。另外,我的系统上有 spark-3.0.1-bin-hadoop2.7。实际上,我想使用 kafka 主题作为 Spark Structured Streaming with python 的来源。我的代码是这样的:

spark = SparkSession \
    .builder \
    .appName("APP") \
    .getOrCreate()

df = spark \
    .readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "sparktest") \
    .option("startingOffsets", "earliest") \
    .load()

df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
df.printSchema()

我通过 spark-submit 使用以下命令运行上述代码:

./spark-submit --packages org.apache.spark:spark-streaming-kafka-0-10_2.12:3.0.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.0.1 /home/spark/PycharmProjects/testSparkStream/KafkaToSpark.py 

代码运行没有任何异常,我在 Spark 站点中收到此输出:

   root
    |-- key: binary (nullable = true)
    |-- value: binary (nullable = true)
    |-- topic: string (nullable = true)
    |-- partition: integer (nullable = true)
    |-- offset: long (nullable = true)
    |-- timestamp: timestamp (nullable = true)
    |-- timestampType: integer (nullable = true)

我的问题是kafka主题充满了数据;但是在输出中运行代码的结果是没有任何数据。请您指导我这里出了什么问题?

【问题讨论】:

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


    【解决方案1】:

    原样的代码不会打印出任何数据,而只会为您提供一次架构。

    您可以按照通用Structured Streaming GuideStructured Streaming + Kafka integration Guide 中给出的说明来查看如何将数据打印到控制台。请记住,在 Spark 中读取数据是一种惰性操作,如果没有操作(通常是 writeStream 操作),则什么都做不了。

    如果您补充如下代码,您应该会看到所选数据(键和值)打印到控制台:

    spark = SparkSession \
              .builder \
              .appName("APP") \
              .getOrCreate()
    
    df = spark\
          .readStream \
          .format("kafka") \
          .option("kafka.bootstrap.servers", "localhost:9092") \
          .option("subscribe", "sparktest") \
          .option("startingOffsets", "earliest") \
          .load()
          
    
    query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
        .writeStream \
        .format("console") \
        .option("checkpointLocation", "path/to/HDFS/dir") \
        .start()
    
    query.awaitTermination()
    

    【讨论】:

      猜你喜欢
      • 2021-12-05
      • 2018-09-06
      • 1970-01-01
      • 2023-03-05
      • 2020-12-30
      • 2020-09-13
      • 2021-04-03
      • 2021-05-08
      • 2019-11-20
      相关资源
      最近更新 更多