Spark读取Kafka数据:交互式模式下的异常行为排查
我在Google Compute Engine虚拟机上部署了单节点Kafka(版本3.5.0,搭配Zookeeper),尝试通过运行在Google Dataproc集群上的简单Spark结构化流代码读取JSON数据。流作业仅负责解析数据并输出到控制台(或目标存储)。
Spark核心处理代码如下:
# 注:schema包含所有输入字段且数据类型正确,测试时也尝试过将所有字段设为StringType() result_df = kafka_stream_df\ .select(from_json(col("value").cast("string"),schema).alias("parsed_value"), col("timestamp").alias("kafka_timestamp"))\ .selectExpr("parsed_value.*","kafka_timestamp")\ .withColumn("event_date",col("timestamp").cast('date'))
我创建了一个单分区的Kafka主题,发布到Kafka的数据为JSON格式。测试时,我分别通过以下两种方式每次发送1条记录:
场景1:通过文件发布数据
使用如下命令,无论数据量大小,Spark流代码都能读取完整数据。我尝试过发布超过5KB的数据,运行正常,可读取所有数据并生成结果。
# 将场景2中使用的相同记录复制到test_data.json文件中 ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic testtopic < /home/user_name/test_data.json
场景2:交互式Shell发布数据
通过VM上的Shell交互式向Kafka主题发布数据,仅当记录大小小于4KB时,我才能读取所有字段。若发布的记录超过4KB,Spark输出的所有字段均为null/空值。
# 在VM上交互式输入场景1中使用的相同记录 ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic testtopic
现象总结:通过文件向Kafka发布数据时,无论记录大小,Spark都能生成正确结果;通过交互式方式发布数据时,仅当记录大小小于4KB时结果正确,超过4KB则所有字段为null/空值。
请问:为何交互式发送超过4KB的消息到Kafka时,Spark会解析出null值,而通过文件导入相同大小的消息时却能正常解析?
核心原因:kafka-console-producer交互式模式的输入限制
kafka-console-producer的交互式模式依赖JLine库处理终端输入,该库默认有4KB的输入缓冲区限制。当你交互式输入超过4KB的JSON数据时,超出缓冲区的内容会被直接截断,导致发送到Kafka的是不完整的JSON字符串。Spark的from_json函数遇到非法JSON时,会返回全null的解析结果。
而通过文件重定向发送数据时,kafka-console-producer会读取整个文件内容作为单条消息(只要文件内没有额外换行符分割消息),不会触发JLine的缓冲区截断逻辑,因此完整的JSON能被正常解析。
解决办法
- 调整交互式生产者的输入缓冲区大小
启动kafka-console-producer时,通过JLine的系统参数修改缓冲区容量,例如设置为10KB:
./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic testtopic \ -Dorg.jline.bufsize=10240
- 优先使用文件输入方式
如果需要发送大消息,建议继续采用文件重定向的方式,避免终端输入的限制。 - 自定义生产者客户端
若需频繁交互式发送大消息,可以自行编写Python/Java等语言的简易生产者客户端,绕过kafka-console-producer的终端限制。
内容的提问来源于stack exchange,提问作者user16798185

