Spark Structured Streaming写入Kafka无报错但无数据问题求助
问题排查步骤
- 缺失
awaitTermination()调用导致任务提前终止
你提供的Kafka写入代码仅调用了.start()启动流任务,没有追加.awaitTermination()阻塞主线程,主线程执行完start逻辑后会直接退出,流任务也会随之终止,没有足够时间处理文件数据并写入Kafka。而console测试代码最后追加了该方法,所以可以正常运行输出。
修复方案:在Kafka写入的.start()后追加.awaitTermination(),代码如下:sdf \ .selectExpr("CAST(country AS STRING) AS key", "to_json(struct(*)) AS value") \ .writeStream \ .format("kafka") \ .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) \ .option("topic", TOPIC) \ .option("checkpointLocation", "/tmp/demo") \ .trigger(processingTime='1 seconds') \ .start() \ .awaitTermination() - 未指定输出模式,默认append模式延迟输出聚合结果
console测试代码显式指定了outputMode("update"),该模式下聚合结果有更新就会实时输出,所以可以马上看到数据。而Kafka写入代码未指定输出模式,默认使用append模式,你配置了5秒的watermark,append模式下聚合结果需要等水印超过对应时间窗口阈值才会输出,短时间内不会生成可写入的数据。
修复方案:Kafka写入代码中追加.outputMode("update")即可和console逻辑一致,实时输出更新后的聚合结果。 - checkpoint目录残留旧元数据
如果你之前使用过/tmp/demo作为checkpoint路径运行过其他流任务,残留的元数据会导致新任务跳过已经处理过的偏移量,不会读取新移入的Parquet文件。
修复方案:删除/tmp/demo目录后重新运行任务。 - Kafka连接/权限问题
无报错不代表Kafka连接正常,部分版本的Spark Kafka Sink会将连接失败、ACL鉴权失败等异常打印在DEBUG级别的日志中,不会抛出到主程序日志。
排查方案:- 确认
KAFKA_BOOTSTRAP_SERVERS地址正确,Spark运行节点可以正常访问Kafka的服务端口 - 如果Kafka开启了ACL鉴权,需要在写入配置中补充对应的sasl认证参数
- 调整Spark日志级别到DEBUG,查找KafkaProducer相关的报错信息
- 确认
- 数据格式不符合Kafka Sink要求
Kafka Sink要求写入的key和value字段必须为String或者Binary类型,且如果字段为null可能被部分Kafka服务端配置拒收。
排查方案:先将selectExpr处理后的DataFrame用console输出,确认key、value两个字段存在、类型正确且非空。
内容的提问来源于stack exchange,提问作者SASM
相关产品推荐
相关产品推荐

