Spark 3.5.1中maxOffsetsPerTrigger参数生效范围及行为确认
关于Spark Structured Streaming
maxOffsetsPerTrigger 参数行为的解答 在Spark 3.5.1版本中,maxOffsetsPerTrigger 参数是按Kafka分区生效的,而非全局限制,你遇到的每个分区读取1条、总计5条消息的情况属于预期行为。
具体来说,该参数的作用是限制每个Kafka分区在单个微批触发器周期内最多读取的偏移量数量。当你将其设置为1时,订阅的5个Kafka分区会各自读取1条未消费的消息,最终微批总计处理5条消息,这完全符合参数的设计逻辑。
如果需要实现全局级别的读取数量限制,你需要在读取数据后额外添加过滤逻辑,比如通过limit算子来限制全局处理的消息总量。
内容的提问来源于stack exchange,提问作者mt_leo
相关产品推荐
相关产品推荐

