Apache Beam中PubSubIO写入是否支持setDelayThreshold()及批量超时设置?
Apache Beam PubSubIO 批量写入相关问题解答
是否存在setDelayThreshold()命令?
Apache Beam的PubSubIO写入组件没有setDelayThreshold()这个方法。
批量写入是否支持时间限制?
是的,除了maxBatchSize()设置批量大小上限外,PubSubIO支持通过setMaxBatchDuration()(Java)或max_batch_duration_secs(Python)配置批量等待的最长时间。当到达这个时间阈值时,无论当前批量消息数量是否达到maxBatchSize()的上限,都会触发批量写入操作。
未达批量上限且管道未结束时,消息会滞留吗?
如果仅配置了maxBatchSize()而未设置时间限制,当消息流入量不足且管道持续运行时,未达到批量大小的消息会滞留在接收器中,直到满足批量大小条件或管道终止。
如何设置批量写入超时时间?
通过同时配置批量大小和时间阈值,即可避免消息长时间滞留,以下是不同语言的配置示例:
Java 示例
PubSubIO.writeStrings() .to("projects/my-project/topics/my-topic") .withMaxBatchSize(1000) // 批量大小上限 .withMaxBatchDuration(Duration.standardSeconds(10)); // 超时时间10秒
Python 示例
beam.io.WriteToPubSub( topic="projects/my-project/topics/my-topic", max_batch_size=1000, # 批量大小上限 max_batch_duration_secs=10 # 超时时间10秒 )
上述配置会让接收器在“收集到1000条消息”或“等待满10秒”两个条件中满足任一条件时,立即执行批量写入。
内容的提问来源于stack exchange,提问作者DiFalco
相关产品推荐
相关产品推荐

