From 882e94f6f456be52c62aa74b93606bb8e7293e2c Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 12 Dec 2014 11:19:44 -0500 Subject: [PATCH] AMQP-451: Support Flushing Batch JIRA: https://jira.spring.io/browse/AMQP-451 - Add `flush()` to support flushing any partial batch. - Implement `Lifecycle` and invoke `flush()` when stopped. --- .../rabbit/core/BatchingRabbitTemplate.java | 24 ++++++++++++++++++- 1 file changed, 23 insertions(+), 1 deletion(-) diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplate.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplate.java index 5f82d2fb..1b42cd82 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplate.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/core/BatchingRabbitTemplate.java @@ -23,6 +23,7 @@ import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.core.support.BatchingStrategy; import org.springframework.amqp.rabbit.core.support.MessageBatch; import org.springframework.amqp.rabbit.support.CorrelationData; +import org.springframework.context.Lifecycle; import org.springframework.scheduling.TaskScheduler; /** @@ -38,7 +39,7 @@ import org.springframework.scheduling.TaskScheduler; * @since 1.4.1 * */ -public class BatchingRabbitTemplate extends RabbitTemplate { +public class BatchingRabbitTemplate extends RabbitTemplate implements Lifecycle { private final BatchingStrategy batchingStrategy; @@ -85,6 +86,13 @@ public class BatchingRabbitTemplate extends RabbitTemplate { } } + /** + * Flush any partial in-progress batches. + */ + public void flush() { + releaseBatches(); + } + private synchronized void releaseBatches() { MessageBatch batch; while ((batch = this.batchingStrategy.releaseBatch()) != null) { @@ -92,4 +100,18 @@ public class BatchingRabbitTemplate extends RabbitTemplate { } } + @Override + public void start() { + } + + @Override + public void stop() { + flush(); + } + + @Override + public boolean isRunning() { + return true; + } + }