【问题标题】:Reading avro messages from Kafka in spark streaming/structured streaming在火花流/结构化流中从 Kafka 读取 avro 消息
【发布时间】:2019-09-20 17:28:39
【问题描述】:

我是第一次使用 pyspark。 火花版本:2.3.0 卡夫卡版本:2.2.0

我有一个 kafka 生产者,它以 avro 格式发送嵌套数据,我正在尝试在 pyspark 中的 spark-streaming/结构化流中编写代码,这会将来自 kafka 的 avro 反序列化为数据帧做转换以 parquet 格式将其写入 s3 . 我能够在 spark/scala 中找到 avro 转换器,但尚未添加对 pyspark 的支持。如何在 pyspark 中进行相同的转换。 谢谢。

【问题讨论】:

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


    【解决方案1】:

    就像你提到的,从 Kafka 读取 Avro 消息并通过 pyspark 解析,没有相同的直接库。但是我们可以通过编写小型包装器来读取/解析 Avro 消息,并在您的 pyspark 流代码中将该函数称为 UDF,如下所示。

    参考: Pyspark 2.4.0, read avro from kafka with read stream - Python

    注意:自 Spark 2.4 以来,Avro 是内置但外部的数据源模块。请按照“Apache Avro 数据源指南”的部署部分部署应用程序。

    参考: https://spark-test.github.io/pyspark-coverage-site/pyspark_sql_avro_functions_py.html

    Spark-提交:

    [调整软件包版本以匹配基于 spark/avro 版本的安装]

    /usr/hdp/2.6.1.0-129/spark2/bin/pyspark --packages org.apache.spark:spark-avro_2.11:2.4.3 --conf spark.ui.port=4064
    

    Pyspark 流式传输代码:

    from pyspark.sql import SparkSession
    from pyspark.sql.functions import *
    from pyspark.sql.types import *
    from pyspark.streaming import StreamingContext
    from pyspark.sql.column import Column, _to_java_column
    from pyspark.sql.functions import col, struct
    from pyspark.sql.functions import udf
    import json
    import csv
    import time
    import os
    
    #  Spark Streaming context :
    
    spark = SparkSession.builder.appName('streamingdata').getOrCreate()
    sc = spark.sparkContext
    ssc = StreamingContext(sc, 20)
    
    #  Kafka Topic Details :
    
    KAFKA_TOPIC_NAME_CONS = "topicname"
    KAFKA_OUTPUT_TOPIC_NAME_CONS = "topic_to_hdfs"
    KAFKA_BOOTSTRAP_SERVERS_CONS = 'localhost.com:9093'
    
    #  Creating  readstream DataFrame :
    
    df = spark.readStream \
         .format("kafka") \
         .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS_CONS) \
         .option("subscribe", KAFKA_TOPIC_NAME_CONS) \
         .option("startingOffsets", "latest") \
         .option("failOnDataLoss" ,"false")\
         .option("kafka.security.protocol","SASL_SSL")\
         .option("kafka.client.id" ,"MCI-CIL")\
         .option("kafka.sasl.kerberos.service.name","kafka")\
         .option("kafka.ssl.truststore.location", "/path/kafka_trust.jks") \
         .option("kafka.ssl.truststore.password", "changeit") \
         .option("kafka.sasl.kerberos.keytab","/path/bdpda.headless.keytab") \
         .option("kafka.sasl.kerberos.principal","bdpda") \
         .load()
    
    
    df1 = df.selectExpr( "CAST(value AS STRING)")
    
    df1.registerTempTable("test")
    
    
    # Deserilzing the Avro code function
    
    from pyspark.sql.column import Column, _to_java_column 
    def from_avro(col): 
         jsonFormatSchema = """
                        {
                         "type": "record",
                         "name": "struct",
                         "fields": [
                           {"name": "col1", "type": "long"},
                           {"name": "col2", "type": "string"}
                                    ]
                         }"""
        sc = SparkContext._active_spark_context 
        avro = sc._jvm.org.apache.spark.sql.avro
        f = getattr(getattr(avro, "package$"), "MODULE$").from_avro
        return Column(f(_to_java_column(col), jsonFormatSchema))
    
    
    spark.udf.register("JsonformatterWithPython", from_avro)
    
    squared_udf = udf(from_avro)
    df1 = spark.table("test")
    df2 = df1.select(squared_udf("value"))
    
    #  Declaring the Readstream Schema DataFrame :
    
    df2.coalesce(1).writeStream \
       .format("parquet") \
       .option("checkpointLocation","/path/chk31") \
       .outputMode("append") \
       .start("/path/stream/tgt31")
    
    
    ssc.awaitTermination()
    

    【讨论】:

    • 为什么需要 df5 和 df1 是同一个东西?另外,正如它对您链接的答案的评论,这不适用于 Confluent Schema Registry Avro 数据
    • 感谢您分享您的 cmets 。你说的对 !! df5 不是必需的,并更正了该代码。正如在给定一般 Avro 格式处理代码的情况下,没有提到关于融合注册表的问题
    • 对,但根据我的经验,这很少是 Kafka 中的内容
    • 仅供参考,这段代码,我已经在 Cloudra 环境中使用经典的 Apache Kafka 实现了
    • 好的,好的,那么您没有使用架构注册表? Confluent 不使用除了 Apache Kafka 之外的任何东西......
    猜你喜欢
    • 2017-04-04
    • 2020-03-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-29
    • 2017-03-27
    • 2020-08-18
    相关资源
    最近更新 更多