【问题标题】:Empty column when deserializing avro from apache kafka with pyspark使用 pyspark 从 apache kafka 反序列化 avro 时为空列
【发布时间】:2019-11-07 02:00:24
【问题描述】:

我正在使用 Kafka、Spark 和 jupyter notebook 进行概念验证,但遇到了一个奇怪的问题。我试图读取从 kafka 到 pyspark 的 Avro 记录。我正在使用融合模式注册表来获取模式以反序列化 avro 消息。 在对 spark 数据帧中的 avro 消息进行反序列化后,结果列是空的,没有任何错误。该列应包含数据,因为当转换为字符串时,某些 avro 字段是可读的。

我也尝试在 Scala 中的 spark-shell 上做这件事(没有 jupyter) 我已经尝试过基于 docker 的 spark 以及独立安装的 spark

我按照这个 SO 主题来获取 from_avro 和 to_avro 函数: Pyspark 2.4.0, read avro from kafka with read stream - Python

jars = ["kafka-clients-2.0.0.jar", "spark-avro_2.11-2.4.3.jar", "spark-        
sql-kafka-0-10_2.11-2.4.3.jar"]
jar_paths = ",".join(["/home/jovyan/work/jars/{}".format(jar) for jar in 
jars])

conf = SparkConf()
conf.set("spark.jars", jar_paths)

spark_session = SparkSession \
    .builder \
    .config(conf=conf)\
    .appName("TestStream") \
    .getOrCreate()

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))) 

schema_registry_url = "http://schema-registry.org"
transaction_schema_name = "Transaction"

transaction_schema = requests.get(" 
{}/subjects/{}/versions/latest/schema".format(schema_registry_url, 
transaction_schema_name)).text


raw_df = spark_session.read.format("kafka") \
# SNIP
    .option("subscribe", "transaction") \
    .option("startingOffsets", "earliest").load()
raw_df = raw_df.limit(1000).cache()

extract_df = raw_df.select(
    raw_df["key"].cast("String"),
    from_avro(raw_df["value"], transaction_schema).alias("value")
)

# This shows data and fields
raw_df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)").show(3, truncate=False)

extract_df.show()

值列的内容为空。我预计解码失败后会出现错误,或者数据会在那里。有谁知道这可能是什么原因,或者如何调试它?

+---+-----+
|key|value|
+---+-----+
|...| [[]]|
|...| [[]]|
|...| [[]]|
|...| [[]]|

【问题讨论】:

  • 不幸的是,spark-avro 不支持 Confluent 的序列化程序写入数据的格式,因此它会失败(通过返回 null/空值)
  • 看看这是否有帮助stackoverflow.com/a/55786881/2308683

标签: apache-spark pyspark apache-kafka avro confluent-schema-registry


【解决方案1】:

您必须手动反序列化数据。截至撰写本文时,PySpark 尚未正式支持 Confluent 模式注册表。您需要使用 Confluent 提供的 KafkaAvroDeSerializer 或 ABRiS(一个 3rd-party Spark avro 库)。

ABRiS:https://github.com/AbsaOSS/ABRiS#using-abris-with-python-and-pyspark

KafkaAvroDeSerializer:Integrating Spark Structured Streaming with the Confluent Schema Registry

原因:Confluent 添加了 5 个额外字节,其中 1 个用于魔术字节,4 个用于架构 ID,在 Avro 数据旁边,[魔术字节|架构 ID|avro 数据],不是典型的 avro 格式。所以需要手动反序列化。

(对不起,我无法发表评论。)

【讨论】:

    猜你喜欢
    • 2019-07-12
    • 2020-07-06
    • 2019-12-04
    • 2019-09-18
    • 1970-01-01
    • 2018-08-11
    • 2019-07-30
    • 2015-08-01
    相关资源
    最近更新 更多