Spark处理Kafka记录时如何将Bytes转换为String?
问题解决:Spark读取Kafka数据时key/value未正确转换为String
你的核心问题是Spark的Dataset是不可变对象,调用selectExpr转换后没有将新生成的Dataset赋值给变量,导致后续输出的仍是原始的字节格式数据。
错误原因
你当前代码里的这一行只是执行了转换操作,但没有保存转换后的结果:
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)");
原df变量依然指向未转换的原始Dataset,所以后续的writeStream输出的还是包含key(二进制)、value(二进制)、topic、partition等字段的原始数据。
修正后的代码
将selectExpr的结果重新赋值给df(或新变量),确保后续流查询使用的是转换后的Dataset:
package org.example; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.streaming.StreamingQuery; import org.apache.spark.sql.streaming.StreamingQueryException; import java.util.concurrent.TimeoutException; public class Kafka { public static void main(String args[]){ SparkSession spark = SparkSession .builder() .appName("Spark-Kafka-Integration") .config("spark.master", "local") .getOrCreate(); Dataset<Row> df = spark .readStream() .format("kafka") .option("kafka.bootstrap.servers", "localhost:9092") .option("subscribe", "quickstart-events") .load(); // 关键:将转换后的结果重新赋值给df df = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)"); StreamingQuery query = null; try { query = df.writeStream() .outputMode("append") .format("console") .start(); } catch (TimeoutException e) { throw new RuntimeException(e); } try { query.awaitTermination(); } catch (StreamingQueryException e) { throw new RuntimeException(e); } } }
额外说明
- 如果需要保留原始的其他字段(比如
topic、timestamp),可以在selectExpr中一并指定,例如:df = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "topic", "timestamp"); - 启动Kafka生产者发送消息后,控制台会输出转换后的字符串格式的key和value(如果发送时指定了key,否则key为null)。
内容的提问来源于stack exchange,提问作者Anurag
相关产品推荐
相关产品推荐

