【问题标题】:Spark Batch Avro Deserialization: Malformed data. Length is negativeSpark Ba​​tch Avro 反序列化:格式错误的数据。长度为负
【发布时间】:2021-10-29 12:39:14
【问题描述】:

我正在通过 Spark 对 Kafka 进行一些批处理。与 Avro 一样序列化的记录。我正在尝试使用消息本身中的确切模式来反序列化该值,但我得到了一个格式错误的记录异常。这是我的代码:

    Dataset<Row> load = sparkSession
            .read()
            .format("kafka")
            .option("kafka.bootstrap.servers", (String) kafkaConfiguration.consumerProperties().get("bootstrap.servers"))
            .option("subscribe", kafkaConfiguration.topicsAsCSV(","))
            .load();

    var schema = new String(Files.readAllBytes(Paths.get("schema.avsc")));

    load.select(from_avro(col("value"), schema)).write().json("/tmp/spark/json");

请注意,架构是从记录的值本身复制而来的。

但是,我得到以下异常:

Caused by: org.apache.avro.AvroRuntimeException: Malformed data. Length is negative: -20
    at org.apache.avro.io.BinaryDecoder.doReadBytes(BinaryDecoder.java:336)
    at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:263)
    at org.apache.avro.io.BinaryDecoder.readString(BinaryDecoder.java:272)
    at org.apache.avro.io.ResolvingDecoder.readString(ResolvingDecoder.java:214)

此错误背后的原因是什么,我该如何解决?谢谢!

【问题讨论】:

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


    【解决方案1】:

    我需要更多信息。 但是,当我尝试使用来自 Confluent 的 Kafka 的 avro 消息时,我遇到了类似的问题。它生成一条 5 字节的消息来存储来自 SchemaRegistry 服务的 SchemaRegitryID 以及它在消息其余部分中的 avro 内容。

    # python example to slice bytes
    schema_id = message[:5]
    content = message[5:]
    

    https://www.confluent.io/blog/consume-avro-data-from-kafka-topics-and-secured-schema-registry-with-databricks-confluent-cloud-on-azure/#parse-bytes 中的更多信息,pyspark 安装了 Java/Scala(抱歉)。

    【讨论】:

    • 我实际上也尝试过跳过前五个字节,但无济于事:/
    猜你喜欢
    • 2017-07-08
    • 2023-04-07
    • 2018-07-20
    • 2020-03-04
    • 2018-07-01
    • 2019-10-20
    • 1970-01-01
    • 2019-01-27
    • 2015-08-01
    相关资源
    最近更新 更多