Spark Kafka API实现Databricks到Event Hubs流写无数据流入问题
问题原因
你的代码没有数据写入是三个核心逻辑错误导致的:
- Spark Structured Streaming的文件流源默认只会消费启动后新放入监听路径的文件,路径下提前存放的历史文件只会在第一次启动流时被扫描,扫描完成后对应文件的偏移量会持久化到checkpoint中,后续重启流时不会重复消费这些已记录的文件。
- 你循环启停同一个checkpoint路径下流查询的写法不符合Structured Streaming的运行逻辑。同一个checkpoint路径唯一对应一个流查询的消费状态,你每次循环启动查询都会读取上次记录的偏移量,没有新文件的情况下查询启动后没有待处理数据,直接空跑结束,所以只会打印日志没有实际写入。
- Kafka兼容的Event Hubs Sink要求输入DataFrame必须包含名为
value的字段(消息体,支持字符串、二进制类型),可选key、headers、partition等字段。你直接读取CSV生成的DataFrame只有自定义业务字段,没有匹配Sink的字段要求,即使读到数据也会写入失败。
对应解决方案
根据你的业务场景二选一即可:
场景1:需要循环重发同一份CSV数据做测试(连续数据生成器)
这种场景不需要用流读CSV,直接将CSV读为批DataFrame循环写入即可,不需要固定checkpoint路径:
from pyspark.sql.functions import to_json, struct import time # 读为批DataFrame,提前转换为Kafka Sink要求的格式 batch_df = ( spark.read.format("csv") .option("header", "true") .schema(location_schema) .load(f"{path}") # 将整行数据转换为JSON格式作为消息value .withColumn("value", to_json(struct(*location_schema.fieldNames()))) .select("value") ) while True: # 直接批写入即可,不需要流启动逻辑 ( batch_df.write .format("kafka") .option("topic", topic) .option("kafka.bootstrap.servers", bootstrap_servers) .option("kafka.sasl.mechanism", "PLAIN") .option("kafka.security.protocol", "SASL_SSL") .option("kafka.sasl.jaas.config", sasl_jaas_config) .save() ) print(f"Wrote {batch_df.count()} records once") time.sleep(5)
场景2:需要持续处理路径下新增的CSV文件
这种场景不需要手动写while循环启停流,Structured Streaming本身自带定时调度能力,启动一次即可持续运行:
from pyspark.sql.functions import to_json, struct # 构建流DataFrame时提前转换字段格式 streaming_df = ( spark.readStream.format("csv") .option("header", "true") .schema(location_schema) .load(f"{path}") .withColumn("value", to_json(struct(*location_schema.fieldNames()))) .select("value") ) # 启动流,设置5秒间隔检查新文件 query = ( streaming_df.writeStream.format("kafka") .option("topic", topic) .option("kafka.bootstrap.servers", bootstrap_servers) .option("kafka.sasl.mechanism", "PLAIN") .option("kafka.security.protocol", "SASL_SSL") .option("kafka.sasl.jaas.config", sasl_jaas_config) .option("checkpointLocation", "/checkpoint") .trigger(processingTime="5 seconds") .start() ) query.awaitTermination()
注意:改用这个方案前先清空原有
/checkpoint路径下的所有文件,旧的偏移量记录会导致历史文件不会被重新消费。
快速排查技巧
对接Event Hubs前先把流输出到控制台验证读逻辑正常,避免Sink配置问题干扰排查:
debug_query = ( streaming_df.writeStream .format("console") .trigger(once=True) .start() ) debug_query.awaitTermination()
控制台能正常打印出value字段内容后,再切换为Kafka格式写入即可。
内容的提问来源于stack exchange,提问作者user14681827
相关产品推荐
相关产品推荐

