【问题标题】:Pyspark 2.4.3, Read Avro format message from Kafka - Pyspark Structured streamingPyspark 2.4.3,从 Kafka 读取 Avro 格式消息 - Pyspark 结构化流
【发布时间】:2020-03-17 21:59:24
【问题描述】:

我正在尝试使用 PySpark 2.4.3 从 Kafka 读取 Avro 消息。基于下面的堆栈溢出链接,我能够转换为 Avro 格式(to_avro)并且代码按预期工作。但 from_avro 无法正常工作并出现以下问题。是否有其他模块支持读取从 Kafka 流式传输的 avro 消息?这是 Cloudra 分发环境。 请就此提出建议。

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

环境详情:

火花:

 / __/__  ___ _____/ /__
    _\ \/ _ \/ _ `/ __/  '_/
   /__ / .__/\_,_/_/ /_/\_\   version 2.1.1.2.6.1.0-129
      /_/

Using Python version 3.6.1 (default, Jul 24 2019 04:52:09)

派斯帕克:

pyspark 2.4.3

Spark_submit:

/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

to_avro

from pyspark.sql.column import Column, _to_java_column 

def from_avro(col, jsonFormatSchema): 
    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)) 


def to_avro(col): 
    sc = SparkContext._active_spark_context 
    avro = sc._jvm.org.apache.spark.sql.avro
    f = getattr(getattr(avro, "package$"), "MODULE$").to_avro
    return Column(f(_to_java_column(col))) 
from pyspark.sql.functions import col, struct


avro_type_struct = """
{
  "type": "record",
  "name": "struct",
  "fields": [
    {"name": "col1", "type": "long"},
    {"name": "col2", "type": "string"}
  ]
}"""


df = spark.range(10).select(struct(
    col("id"),
    col("id").cast("string").alias("id2")
).alias("struct"))
avro_struct_df = df.select(to_avro(col("struct")).alias("avro"))
avro_struct_df.show(3)
+----------+
|      avro|
+----------+
|[00 02 30]|
|[02 02 31]|
|[04 02 32]|
+----------+
only showing top 3 rows

来自_avro:

avro_struct_df.select(from_avro("avro", avro_type_struct)).show(3)

错误信息:

Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/usr/hdp/2.6.1.0-129/spark2/python/pyspark/sql/dataframe.py", line 993, in select
    jdf = self._jdf.select(self._jcols(*cols))
  File "/usr/hdp/2.6.1.0-129/spark2/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
  File "/usr/hdp/2.6.1.0-129/spark2/python/pyspark/sql/utils.py", line 63, in deco
    return f(*a, **kw)
  File "/usr/hdp/2.6.1.0-129/spark2/python/lib/py4j-0.10.4-src.zip/py4j/protocol.py", line 319, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o61.select.
: java.lang.NoSuchMethodError: org.apache.avro.Schema.getLogicalType()Lorg/apache/avro/LogicalType;
        at org.apache.spark.sql.avro.SchemaConverters$.toSqlTypeHelper(SchemaConverters.scala:66)
        at org.apache.spark.sql.avro.SchemaConverters$$anonfun$1.apply(SchemaConverters.scala:82)

【问题讨论】:

    标签: apache-spark pyspark avro spark-avro


    【解决方案1】:

    您的 Spark 版本实际上是 2.1.1,因此您不能使用 Spark 中包含的 spark-avro 包的 2.4.3 版本

    您需要使用 Databricks 中的那个

    是否还有其他模块支持读取从 Kafka 流式传输的 avro 消息?

    你可以使用普通的 kafka Python 库,而不是 Spark

    【讨论】:

      【解决方案2】:

      Spark 2.4.0 支持 to_avrofrom_avro 函数,但仅适用于 Scala and Java。那么你的方法应该没问题,只要使用适当的spark versionspark-avro 包。

      在使用 Spark Structure Streaming 来使用 Kafka 消息时,我更喜欢另一种方法是使用带有 fastavro python 库的 UDF。 fastavro 相对较快,因为它使用了 C 扩展。我已经将它用于我们的生产几个月了,没有任何问题。

      如下面代码sn-p所示,Kafka主消息承载在kafka_dfvalues列中。出于演示目的,我使用了一个简单的 avro 模式,其中包含 2 列 col1col2deserialize_avro UDF 函数的返回是一个元组,对应于 avro 模式中描述的字段数。然后将流写入控制台以进行调试。

      from pyspark.sql import SparkSession
      import pyspark.sql.functions as psf
      from pyspark.sql.types import *
      import io
      import fastavro
      
      def deserialize_avro(serialized_msg):
          bytes_io = io.BytesIO(serialized_msg)
          bytes_io.seek(0)
          avro_schema = {
                          "type": "record",
                          "name": "struct",
                          "fields": [
                            {"name": "col1", "type": "long"},
                            {"name": "col2", "type": "string"}
                          ]
                        }
      
          deserialized_msg = fastavro.schemaless_reader(bytes_io, avro_schema)
      
          return (    deserialized_msg["col1"],
                      deserialized_msg["col2"]
                  )
      
      if __name__=="__main__":
        spark = SparkSession \
              .builder \
              .appName("consume kafka message") \
              .getOrCreate()
      
        kafka_df = spark \
                    .readStream \
                    .format("kafka") \
                    .option("kafka.bootstrap.servers", "kafka01-broker:9092") \
                    .option("subscribe", "topic_name") \
                    .option("stopGracefullyOnShutdown", "true") \
                    .load()
      
        df_schema = StructType([
                    StructField("col1", LongType(), True),
                    StructField("col2", StringType(), True)
                ])
      
        avro_deserialize_udf = psf.udf(deserialize_avro, returnType=df_schema)
        parsed_df = kafka_df.withColumn("avro", avro_deserialize_udf(psf.col("value"))).select("avro.*")
      
        query = parsed_df.writeStream.format("console").option("truncate", "true").start()
        query.awaitTermination()
      

      【讨论】:

        猜你喜欢
        • 2019-09-20
        • 2017-04-04
        • 2019-07-08
        • 2019-09-28
        • 1970-01-01
        • 1970-01-01
        • 2019-01-29
        • 2023-03-25
        • 1970-01-01
        相关资源
        最近更新 更多