【问题标题】:Print Pyspark Streaming Data from a kafka Topic从 kafka 主题打印 Pyspark 流数据
【发布时间】:2021-07-30 06:47:51
【问题描述】:

我是 kafka 和 pyspark 的新手,正在尝试编写简单的程序,所以我有 2 个 JSon 格式的 kafka Topics 文件,我正在从 pyspark 流中读取。

我的Producer代码如下:

  from kafka import *
import json
import time
import boto3
import json
from Consumer_Group import *
from json import loads
class producer :
            def json_serializer(data):
                    return json.dumps(data).encode("utf-8")

            def read_s3():
                p1 = KafkaProducer(bootstrap_servers=['localhost:9092'], value_serializer=producer.json_serializer)
                s3 = boto3.resource('s3')
                bucket = s3.Bucket('kakfa')
                for obj in bucket.objects.all():
                    key = obj.key
                    body = obj.get()['Body'].read().decode('utf-8')
                p1.send("Uber_Eats",body)
                p1.flush()

我的消费者代码如下:

from pyspark.sql import SparkSession
from kafka import *
import time
class consumer:
                def read_from_topic(self,spark):
                        df = spark.readStream \
                            .format("kafka") \
                            .option("kafka.bootstrap.servers", "localhost:9092") \
                            .option("subscribe", "Uber_Eats") \
                             .option("startingOffsets", "earliest") \
                            .load()
                        df.createOrReplaceTempView("kafka")
                        spark.sql("select * from kafka")
                        print(df.isStreaming())
                                  


                def get_consumer(self):
                    consumer = KafkaConsumer("Uber_Eats", group_id='group1', bootstrap_servers=
                    "localhost:9092")
                    return  consumer

                def print_details(self,c1):
                    #    self.consumer=self.get_consumer(self)
                        # Read and print message from consumer
                     try:
                                for msg in c1:
                                    print(msg.topic, msg.value)
                                print("Done")
                     except Exception  as e:
                                print(e)

主类:

from Producer_Group import *
from Consumer_Group import *
from Spark_Connection import *
class client:
    def transfer(self):
        spark = connection.get_connection(self)
        producer.read_s3()
        c1 = consumer.get_consumer(spark)
        consumer.read_from_topic(self,spark)
      #  consumer.print_details(self,c1)

c=client()
c.transfer()

我正在读入 kafka 主题的 S3 中的示例数据:

{
    
        {
            "Customer Number": "1",
            "Customer Name": "Aditya",
            "Restaurant Number": "2201",
            "Restaurant NameOrdered": "Bawarchi",
            "Number of Items": "3",
            "price": "10",
            "Operating Start hours": "9:00",
            "Operating End hours": "23:00"
        },
        {
            "Customer Number": "2",
            "Customer Name": "Sarva",
            "Restaurant Number": "2202",
            "Restaurant NameOrdered": "Sarvana Bhavan",
            "Number of Items": "4",
            "price": "20",
            "Operating Start hours": "8:00",
            "Operating End hours": "20:00"
        },
        {
            "Customer Number": "3",
            "Customer Name": "Kala",
            "Restaurant Number": "2203",
            "Restaurant NameOrdered": "Taco Bell",
            "Number of Items": "5",
            "price": "30",
            "Operating Start hours": "11:00",
            "Operating End hours": "21:00"
        }
    
}

到目前为止我尝试了什么::我尝试在控制台上打印以检查条件,如果它通过了,则将其插入数据库。为了检查条件,我正在从“read_from_topic”函数读取数据并创建一个视图(createOrReplaceTempView)来查看数据,但没有打印,有人可以指导我如何打印并验证我的条件或数据是否被正确读取?

提前致谢!!!!

【问题讨论】:

  • 我想我之前已经告诉过你,KafkaConsumerreadStream.format("kafka") 是两个完全独立的库,如果你想使用 Spark,你不应该使用第一个,你不应该如果唯一的目标是从 Kafka 消费,请使用第二个......那么,你能澄清一下你的目标吗?除此之外,您只显示了一个您从未调用过的类定义,那么这段代码应该打印什么?
  • 我已经包含了“主类”和示例数据(JSON)。请看一下所以我正在尝试从 kafka 主题中读取,并且对于每个客户都有“操作开始时间”,如果它大于 8:00,那么我应该插入 mysql 数据库,否则忽略,这是我的要求。
  • 我的目标是只从 kafka 主题中消费,并在 spark 流中进行转换,并基于该插入到数据库中
  • @OneCricketeer 感谢您的帮助
  • 我的意思是,Python 的KafkaConsumer 也可以“进行转换”,您可以将其与其他库结合起来“写入数据库”,所以仍然不清楚为什么您认为需要 Spark?此外,我什至不会为此使用 Python,因为 Kafka 附带 Kafka Connect 正是为了写入外部系统rmoff.net/2021/03/12/… 除此之外,当它运行时会发生什么?因为您实际上应该在consumer.read_from_topic(self,spark) 周围遇到错误

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


【解决方案1】:

创建视图 (createOrReplaceTempView) 以查看数据,但没有打印任何内容

因为spark.sql 返回一个新的Dataframe。

如果你想打印它,那么你需要

spark.sql("select * from kafka").show()

但是,仅此一项就至少是两个字节数组列,而不是 JSON 字符串,因此您需要 define a schema at some point to extract anything or CAST to at least have human readable data

还值得指出的是,您显示的数据不是有效的 JSON,boto3 不是必需的,因为 Spark 可以从 S3 本身读取文件(因此并不严格需要 Kafka,因为您可以只使用 S3将数据直接传输到您的最终位置,中间使用 Spark persist() 函数)

【讨论】:

  • 我正在尝试使用您提供的 .show 方法,但给了我错误。 pyspark.sql.utils.AnalysisException:必须使用 writeStream.start() 执行带有流源的查询;
  • 正确。您需要 start() 和 awaitTermination。这是 Spark 源代码中的一个示例github.com/apache/spark/blob/master/examples/src/main/scala/org/…
  • 是的,我试过了,它给了我一个不同的错误,RROR StreamMetadata: Error writing stream metadata StreamMetadata(68220860-8bb2-4058-8799-64d5ef5fcc7d) to file:/C:/Users /komu0/AppData/Local/Temp/temporary-7975a168-d556-4ec3-abfb-f23c8c8508cd/metadata ExitCodeException exitCode=-1073741515:
  • 我不知道这个。我也没有在windows上开发,但这里似乎有一些解决方案stackoverflow.com/questions/45947375/…
  • 当然,非常感谢
猜你喜欢
  • 2020-05-08
  • 2021-08-11
  • 2018-11-18
  • 2016-03-24
  • 1970-01-01
  • 2021-03-01
  • 2019-04-21
  • 2020-02-27
  • 1970-01-01
相关资源
最近更新 更多