如何在Spark Structured Streaming中控制每个Trigger处理的记录数量
Spark Structured Streaming单触发器处理记录数控制方案
你之前配置的参数不生效是因为参数适用场景不匹配,具体实现方式如下:
- 你已经用到的两个参数的作用说明:
inputRowsPerSecond是背压速率限制参数,仅对Kafka、Socket这类速率感知型数据源生效,控制的是整体每秒摄入速率,不是单触发器的最大记录数,文件类数据源下该参数不会生效。maxFilesPerTrigger是仅文件数据源适用的参数,控制每个触发器最多读取的新文件数量,无法直接限制记录数,单文件记录数波动大的场景下无法满足按记录数控制的需求。
通用控制方案
你可以通过maxOffsetsPerTrigger参数实现单触发器最大记录数控制,该参数对所有结构化流数据源生效,作用就是限制每个触发器处理的最大记录数,直接在writeStream的配置项中添加即可,示例如下:
.option("maxOffsetsPerTrigger", 1000) // 数值按需调整,代表单触发器最多处理1000条记录
该参数优先级高于maxFilesPerTrigger,两个参数同时配置时会同时生效,先触发哪个阈值就按哪个规则停止当前批次的读取。
你的代码修改示例
比如需要每个触发器最多处理100条记录,调整后代码如下:
import org.apache.spark.sql.streaming.Trigger val checkpointPath = "/user/akash-singh.bisht@unilever.com/dbacademy/developer-foundations-capstone/checkpoint/orders" val devicesQuery = df.writeStream .outputMode("append") .format("delta") .queryName("orders") .trigger(Trigger.ProcessingTime("1 second")) // 新增配置:单触发器最大处理100条记录 .option("maxOffsetsPerTrigger", 100) .option("maxFilesPerTrigger", 1) .option("checkpointLocation",checkpointPath) .table("orders")
如果使用Spark 3.3及以上版本,还可以搭配minOffsetsPerTrigger参数配置每个触发器最少处理的记录数,避免频繁触发小批次任务,提升资源利用率。
内容的提问来源于stack exchange,提问作者akash singh Bisht
相关产品推荐
相关产品推荐

