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文件)。
- 配置Record Reader为
- 确保
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里的数据格式:
- 用Kafka控制台消费者查看原始数据:
bin/kafka-console-consumer.sh --bootstrap-server host:port --topic topic_name --from-beginning- 如果看到明文JSON,说明是文本格式;如果是乱码,大概率是Avro二进制。
- 验证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
相关产品推荐
相关产品推荐

