You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Structured Streaming读取Kafka中Avro事件失败求助

问题分析与解决方案

首先,这个Malformed data. Length is negative错误的核心原因很明确:Spark尝试用Avro解码器解析非Avro格式的数据。从你提供的Kafka示例数据来看,Topic里存的是明文JSON字符串,但你的Spark代码却用from_avro()来解析value字段——Avro是二进制序列化格式,用它去读文本JSON自然会解码失败,因为解码器会把JSON的字节当作Avro二进制的长度标识来读取,从而出现负数长度的异常。

接下来分几种情况给出解决办法:

情况1:你确实想发送Avro格式数据到Kafka(符合最初的设计)

问题出在NiFi的流处理环节,你需要确保NiFi正确将JSON转换成Avro二进制再发送到Kafka,而不是直接推送JSON字符串。具体配置步骤:

  • 添加ConvertRecord处理器到NiFi流中,放在JSON数据源和PutKafka之间。
    • 配置Record Reader为JsonTreeReader,确保能正确读取输入的JSON数据。
    • 配置Record Writer为AvroRecordSetWriter,并指定与Spark中schema.avsc完全一致的Avro Schema(可以直接导入该schema文件)。
  • 确保PutKafka发送的是ConvertRecord输出的Avro二进制内容,而不是原始JSON。

情况2:你实际需要消费的是JSON格式数据(调整Spark代码)

如果NiFi推送的就是JSON字符串(或者你不需要Avro序列化),那么需要修改Spark代码,用JSON解析代替Avro解析:

import org.apache.spark.sql.types.{StructField, StructType, StringType}
import org.apache.spark.sql.functions.{col, from_json}

val spark = SparkSession.builder.appName("Spark-Kafka-Integration").master("local").getOrCreate()

// 定义与JSON结构匹配的StructType
val jsonSchema = StructType(Seq(
  StructField("host", StringType),
  StructField("event", StringType),
  StructField("connectiontype", StringType),
  StructField("user", StringType),
  StructField("eventtimestamp", StringType)
))

val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "host:port")
  .option("subscribe", "topic_name")
  .load()

// 先将Kafka的二进制value转为字符串,再用from_json解析
val df1 = df.select(
  from_json(col("value").cast(StringType), jsonSchema).as("data")
).select("data.*")

df1.writeStream
  .format("console")
  .option("truncate","false")
  .start()
  .awaitTermination()

验证数据格式的关键步骤

在调整代码或NiFi流之前,先确认Kafka Topic里的数据格式:

  1. 用Kafka控制台消费者查看原始数据:
    bin/kafka-console-consumer.sh --bootstrap-server host:port --topic topic_name --from-beginning
    
    • 如果看到明文JSON,说明是文本格式;如果是乱码,大概率是Avro二进制。
  2. 验证Avro数据合法性(如果是Avro格式):
    bin/kafka-console-consumer.sh --bootstrap-server host:port --topic topic_name --from-beginning | avro-tools tojson --schema-file schema.avsc
    
    • 如果能正常输出JSON,说明Avro数据格式正确;如果报错,说明序列化过程有问题。

额外注意:版本兼容性

HDP 3.1.0生态下的组件版本要确保兼容:

  • Spark 2.4.0对应的spark-avro包版本要匹配(比如spark-avro_2.11-2.4.0.jar),避免因Avro版本不一致导致的序列化问题。
  • NiFi使用的Avro版本要和Spark保持一致(HDP 3.1默认Avro版本是1.8.2),可以在NiFi的AvroRecordSetWriter配置中确认Schema的兼容性。

内容的提问来源于stack exchange,提问作者Chhaya Vishwakarma

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:46:12