如何使用Spark写入带Kafka消息头的Parquet文件
1. Schema定义正确性与正确添加方式
你当前定义的rfSchema完全无效,写法是错的。
Spark Structured Streaming的Kafka数据源自带固定Schema,不会读取你自定义的StructType配置,只要调用.format("kafka").load(),返回的Dataset固定包含以下列:
key:binary(二进制字节数组)类型,对应Kafka消息键value:binary类型,对应Kafka消息值headers:array<structkey:string,value:binary>类型,对应Kafka消息头数组- 元数据列:
topic(string)、partition(int)、offset(long)、timestamp(timestamp)、timestampType(int)
另外你配置的key.deserializer、value.deserializer两个Kafka原生消费者参数,在Structured Streaming中不会生效,Spark不会调用Kafka自带的反序列化器,读出来的key、value默认永远是二进制类型,需要你自己在后续逻辑中做类型转换。
你不需要在读Kafka阶段定义Schema,正确的做法是读取到Kafka原始二进制数据后,通过SQL表达式做字段提取、类型转换,组装成你需要的业务Schema。
2. Parquet输出结构控制方式
Parquet文件的Schema和你写入前的Dataset Schema完全一致,你现在直接把Kafka原始Dataset写入,输出的自然是Kafka源自带的默认列,和你预期的结构不符。
调整方式非常直接:在写入前通过select/selectExpr/withColumn等算子,把数据裁剪、转换成你需要的(kafkaKey, kafkaHeader1, kafkaHeader2, byteArr)四列结构,再调用写入逻辑即可。
另外你当前的写入路径写的是home/xxx/spark3/,缺少根路径斜杠,属于相对路径,会写入到程序运行目录下,建议改成绝对路径/home/xxx/spark3/。如果不需要按Kafka分区做目录分区,可以删掉partitionBy("partition")配置,避免额外的分区目录。
3. selectExpr语句作用与必要性
你看到的ds.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "headers")语句作用很简单:选择key、value、headers三列,同时把原本二进制类型的key、value强制转换为STRING类型,方便处理文本类型的Kafka消息、或者打印调试。
这个语句完全不是必须添加的,只有当你的key、value本身是UTF-8编码的文本内容时,转成STRING才有意义。你当前场景下key需要转Long类型、value需要保留原始字节数组,完全不需要照搬这段示例代码,按自己的业务字段需求做转换即可。
修正后的核心参考代码
// 读取Kafka原始数据,不需要提前自定义Schema Dataset<Row> kafkaRawDs = spark .readStream() .format("kafka") .option("kafka.bootstrap.servers", "10.0.0.0:30526") .option("subscribe", "test.topic") .option("includeHeaders", "true") .option("max.poll.records", "4000") .option("group.id", "testSpark") .option("startingOffsets", "earliest") .option("failOnDataLoss", "false") .load(); // 字段提取、类型转换,组装成目标Schema Dataset<Row> targetDs = kafkaRawDs .selectExpr( // 二进制key转Long类型,命名为kafkaKey "CAST(key AS LONG) as kafkaKey", // 原始value直接保留为字节数组,命名为byteArr "value as byteArr", // 从headers数组中过滤出key为kafkaHeader1的第一个字段值 "element_at(filter(headers, h -> h.key = 'kafkaHeader1'), 1).value as h1Bin", // 提取kafkaHeader2的字段值 "element_at(filter(headers, h -> h.key = 'kafkaHeader2'), 1).value as h2Bin" ) .selectExpr( "kafkaKey", // 二进制header值转字符串 "CAST(h1Bin AS STRING) as kafkaHeader1", "CAST(h2Bin AS STRING) as kafkaHeader2", "byteArr" ); // 写入Parquet SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd"); String currentDate= format.format(new Date()); targetDs.writeStream() .option("checkpointLocation", "/home/xxx/spark3/streamingCheckpoint") .format("parquet") .outputMode(OutputMode.Append()) .start("/home/xxx/spark3/"+currentDate) // 替代原来的Thread.sleep,流任务会持续运行直到主动终止,更符合流处理逻辑 .awaitTermination();
额外注意点
- 不要用
Thread.sleep控制流任务运行时长,awaitTermination()是Spark官方推荐的方式,避免sleep时长不足导致任务提前退出、或者时长过长浪费资源 - 如果你的Kafka消息头的值不是UTF-8编码,不要直接用
CAST转STRING,需要按实际编码做自定义转码 - 如果业务允许空值,记得把Schema对应字段的nullable属性调整为true,避免空值时报错
内容的提问来源于stack exchange,提问作者code123

