diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java index 597d27e536..86cbd87c2c 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java @@ -70,73 +70,73 @@ import org.springframework.util.ErrorHandler; public class AmqpChannelFactoryBean extends AbstractFactoryBean implements SmartLifecycle, BeanNameAware { - private volatile AbstractAmqpChannel channel; - - private volatile List interceptors; + private final AmqpTemplate amqpTemplate = new RabbitTemplate(); private final boolean messageDriven; - private final AmqpTemplate amqpTemplate = new RabbitTemplate(); + private AbstractAmqpChannel channel; - private volatile AmqpAdmin amqpAdmin; + private List interceptors; - private volatile FanoutExchange exchange; + private AmqpAdmin amqpAdmin; - private volatile String queueName; + private FanoutExchange exchange; - private volatile boolean autoStartup = true; + private String queueName; - private volatile Advice[] adviceChain; + private boolean autoStartup = true; - private volatile Integer concurrentConsumers; + private Advice[] adviceChain; - private volatile Integer consumersPerQueue; + private Integer concurrentConsumers; - private volatile ConnectionFactory connectionFactory; + private Integer consumersPerQueue; - private volatile MessagePropertiesConverter messagePropertiesConverter; + private ConnectionFactory connectionFactory; - private volatile ErrorHandler errorHandler; + private MessagePropertiesConverter messagePropertiesConverter; - private volatile Boolean exposeListenerChannel; + private ErrorHandler errorHandler; - private volatile Integer phase; + private Boolean exposeListenerChannel; - private volatile Integer prefetchCount; + private Integer phase; - private volatile boolean isPubSub; + private Integer prefetchCount; - private volatile Long receiveTimeout; + private boolean isPubSub; - private volatile Long recoveryInterval; + private Long receiveTimeout; - private volatile Long shutdownTimeout; + private Long recoveryInterval; - private volatile String beanName; + private Long shutdownTimeout; - private volatile AcknowledgeMode acknowledgeMode; + private String beanName; - private volatile boolean channelTransacted; + private AcknowledgeMode acknowledgeMode; - private volatile Executor taskExecutor; + private boolean channelTransacted; - private volatile PlatformTransactionManager transactionManager; + private Executor taskExecutor; - private volatile TransactionAttribute transactionAttribute; + private PlatformTransactionManager transactionManager; - private volatile Integer txSize; + private TransactionAttribute transactionAttribute; - private volatile Integer maxSubscribers; + private Integer batchSize; - private volatile Boolean missingQueuesFatal; + private Integer maxSubscribers; - private volatile MessageDeliveryMode defaultDeliveryMode; + private Boolean missingQueuesFatal; - private volatile Boolean extractPayload; + private MessageDeliveryMode defaultDeliveryMode; - private volatile AmqpHeaderMapper outboundHeaderMapper = DefaultAmqpHeaderMapper.outboundMapper(); + private Boolean extractPayload; - private volatile AmqpHeaderMapper inboundHeaderMapper = DefaultAmqpHeaderMapper.inboundMapper(); + private AmqpHeaderMapper outboundHeaderMapper = DefaultAmqpHeaderMapper.outboundMapper(); + + private AmqpHeaderMapper inboundHeaderMapper = DefaultAmqpHeaderMapper.inboundMapper(); private boolean headersLast; @@ -314,8 +314,18 @@ public class AmqpChannelFactoryBean extends AbstractFactoryBean> extends AmqpPollableMessageChannelSpec { - private final List adviceChain = new LinkedList(); + private final List adviceChain = new LinkedList<>(); AmqpMessageChannelSpec(ConnectionFactory connectionFactory) { super(new AmqpChannelFactoryBean(true), connectionFactory); @@ -205,16 +206,28 @@ public class AmqpMessageChannelSpec> extends * Configure the txSize. * @param txSize the txSize. * @return the spec. - * @see org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer#setTxSize(int) + * @deprecated since 5.2 in favor of {@link #batchSize(int)} */ + @Deprecated public S txSize(int txSize) { - this.amqpChannelFactoryBean.setTxSize(txSize); + return batchSize(txSize); + } + + /** + * Configure the batch size. + * @param batchSize the batchSize. + * @return the spec. + * @since 5.2 + * @see org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer#setBatchSize(int) + */ + public S batchSize(int batchSize) { + this.amqpChannelFactoryBean.setBatchSize(batchSize); return _this(); } @Override protected AbstractAmqpChannel doGet() { - this.amqpChannelFactoryBean.setAdviceChain(this.adviceChain.toArray(new Advice[this.adviceChain.size()])); + this.amqpChannelFactoryBean.setAdviceChain(this.adviceChain.toArray(new Advice[0])); return super.doGet(); }