使用sparklyr连接Kafka无法读取topic数据 同参数Python可正常连通
问题根因
- 无数据:Spark Structured Streaming 读取Kafka默认从最新偏移量(
latest)开始消费,你Topic中已有的800条都是历史存量数据,自然不会被拉取,返回的只有Kafka数据源的固定表结构,没有实际行数据。 - 延迟高:延迟是Spark Kafka连接器默认的内部超时配置导致的,在没有读到新数据的情况下,会等到内部轮询超时才会返回空结果,和你直接调用Kafka客户端的逻辑不同。
修正方案
- 补充偏移量重置配置
在stream_read_kafka的options里添加kafka.auto.offset.reset = "earliest",和你直接调用KafkaConsumer的配置对齐,从最早的偏移量开始拉取存量数据。 - 调整读取逻辑适配测试场景
如果你只是需要读取存量数据做开发测试,建议优先用批处理方式读取Kafka,避免流式查询的持续监听逻辑导致的卡住问题;如果需要保留流式处理能力,可以配置触发条件和超时参数。 - 确认版本兼容
你当前使用的org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.0和Spark 2.4.7、Kafka 2.8版本兼容,不需要更换依赖包。
可正常运行的参考代码
library(sparklyr) sc <- spark_connect(master = "local", version = "2.4", config = list( sparklyr.shell.packages = "org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.0" )) # 批处理方式读取存量Kafka数据,无延迟 Z <- spark_read_kafka( sc, options = list( kafka.bootstrap.servers="11.22.33.44:5555", subscribe = "mytopic", kafka.auto.offset.reset = "earliest", # 可选:指定偏移量范围,控制读取的数据量 startingOffsets = "earliest", endingOffsets = "latest" ) ) # 直接拉取数据 collect(Z) # 如果需要保留流式处理能力,用以下写法 stream_Z <- stream_read_kafka( sc, options = list( kafka.bootstrap.servers="11.22.33.44:5555", subscribe = "mytopic", kafka.auto.offset.reset = "earliest" ) ) # 流式处理可以先写转换逻辑,再通过stream_write_*系列函数启动流任务即可,不要直接collect流式表
附加说明
如果后续需要做窗口计算、流关联等有状态操作,只需要在流式读取的结果上写对应的sparklyr转换逻辑,再通过stream_write_*系列函数启动流任务即可,不需要自行实现Kafka消费和状态管理逻辑。
内容的提问来源于stack exchange,提问作者r2evans
相关产品推荐
相关产品推荐

