Fix depreactions
This commit is contained in:
@@ -86,12 +86,6 @@ public class KafkaConsumerProperties {
|
||||
|
||||
}
|
||||
|
||||
/**
|
||||
* When true the offset is committed after each record, otherwise the offsets for the complete set of records
|
||||
* received from the poll() are committed after all records have been processed.
|
||||
*/
|
||||
@Deprecated
|
||||
private boolean ackEachRecord;
|
||||
|
||||
/**
|
||||
* When true, topic partitions is automatically rebalanced between the members of a consumer group.
|
||||
@@ -99,14 +93,6 @@ public class KafkaConsumerProperties {
|
||||
*/
|
||||
private boolean autoRebalanceEnabled = true;
|
||||
|
||||
/**
|
||||
* Whether to autocommit offsets when a message has been processed.
|
||||
* If set to false, a header with the key kafka_acknowledgment of the type org.springframework.kafka.support.Acknowledgment header
|
||||
* is present in the inbound message. Applications may use this header for acknowledging messages.
|
||||
*/
|
||||
@Deprecated
|
||||
private boolean autoCommitOffset = true;
|
||||
|
||||
/**
|
||||
* Controlling the container acknowledgement mode. This is the preferred way to control the ack mode on the
|
||||
* container instead of the deprecated autoCommitOffset property.
|
||||
@@ -152,12 +138,6 @@ public class KafkaConsumerProperties {
|
||||
*/
|
||||
private KafkaProducerProperties dlqProducerProperties = new KafkaProducerProperties();
|
||||
|
||||
/**
|
||||
* @deprecated No longer used by the binder.
|
||||
*/
|
||||
@Deprecated
|
||||
private int recoveryInterval = 5000;
|
||||
|
||||
/**
|
||||
* List of trusted packages to provide the header mapper.
|
||||
*/
|
||||
@@ -230,53 +210,6 @@ public class KafkaConsumerProperties {
|
||||
*/
|
||||
private boolean reactiveAtMostOnce;
|
||||
|
||||
/**
|
||||
* @return if each record needs to be acknowledged.
|
||||
*
|
||||
* When true the offset is committed after each record, otherwise the offsets for the complete set of records
|
||||
* received from the poll() are committed after all records have been processed.
|
||||
*
|
||||
* @deprecated since 3.1 in favor of using {@link #ackMode}
|
||||
*/
|
||||
@Deprecated
|
||||
public boolean isAckEachRecord() {
|
||||
return this.ackEachRecord;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param ackEachRecord
|
||||
*
|
||||
* @deprecated in favor of using {@link #ackMode}
|
||||
*/
|
||||
@Deprecated
|
||||
public void setAckEachRecord(boolean ackEachRecord) {
|
||||
this.ackEachRecord = ackEachRecord;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return is autocommit offset enabled
|
||||
*
|
||||
* Whether to autocommit offsets when a message has been processed.
|
||||
* If set to false, a header with the key kafka_acknowledgment of the type org.springframework.kafka.support.Acknowledgment header
|
||||
* is present in the inbound message. Applications may use this header for acknowledging messages.
|
||||
*
|
||||
* @deprecated since 3.1 in favor of using {@link #ackMode}
|
||||
*/
|
||||
@Deprecated
|
||||
public boolean isAutoCommitOffset() {
|
||||
return this.autoCommitOffset;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param autoCommitOffset
|
||||
*
|
||||
* @deprecated in favor of using {@link #ackMode}
|
||||
*/
|
||||
@Deprecated
|
||||
public void setAutoCommitOffset(boolean autoCommitOffset) {
|
||||
this.autoCommitOffset = autoCommitOffset;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return Container's ack mode.
|
||||
*/
|
||||
@@ -348,26 +281,6 @@ public class KafkaConsumerProperties {
|
||||
this.autoCommitOnError = autoCommitOnError;
|
||||
}
|
||||
|
||||
/**
|
||||
* No longer used.
|
||||
* @return the interval.
|
||||
* @deprecated No longer used by the binder
|
||||
*/
|
||||
@Deprecated
|
||||
public int getRecoveryInterval() {
|
||||
return this.recoveryInterval;
|
||||
}
|
||||
|
||||
/**
|
||||
* No longer used.
|
||||
* @param recoveryInterval the interval.
|
||||
* @deprecated No longer needed by the binder
|
||||
*/
|
||||
@Deprecated
|
||||
public void setRecoveryInterval(int recoveryInterval) {
|
||||
this.recoveryInterval = recoveryInterval;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return is auto rebalance enabled
|
||||
*
|
||||
|
||||
@@ -321,18 +321,6 @@ public class KafkaMessageChannelBinder extends
|
||||
this.producerListener = producerListener;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set a {@link ClientFactoryCustomizer} for the {@link ProducerFactory} and {@link ConsumerFactory} created inside
|
||||
* the binder.
|
||||
*
|
||||
* @param customizer the client factory customizer
|
||||
* @deprecated in favor of {@link #addClientFactoryCustomizer(ClientFactoryCustomizer)}.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setClientFactoryCustomizer(ClientFactoryCustomizer customizer) {
|
||||
addClientFactoryCustomizer(customizer);
|
||||
}
|
||||
|
||||
public void addClientFactoryCustomizer(ClientFactoryCustomizer customizer) {
|
||||
if (customizer != null) {
|
||||
this.clientFactoryCustomizers.add(customizer);
|
||||
@@ -686,17 +674,7 @@ public class KafkaMessageChannelBinder extends
|
||||
messageListenerContainer.setBeanName(destination + ".container");
|
||||
// end of these won't be needed...
|
||||
ContainerProperties.AckMode ackMode = extendedConsumerProperties.getExtension().getAckMode();
|
||||
if (ackMode == null) {
|
||||
if (extendedConsumerProperties.getExtension().isAckEachRecord()) {
|
||||
ackMode = ContainerProperties.AckMode.RECORD;
|
||||
}
|
||||
else {
|
||||
if (!extendedConsumerProperties.getExtension().isAutoCommitOffset()) {
|
||||
messageListenerContainer.getContainerProperties()
|
||||
.setAckMode(ContainerProperties.AckMode.MANUAL);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (ackMode != null) {
|
||||
if ((extendedConsumerProperties.isBatchMode() && ackMode != ContainerProperties.AckMode.RECORD) ||
|
||||
!extendedConsumerProperties.isBatchMode()) {
|
||||
@@ -1425,32 +1403,12 @@ public class KafkaMessageChannelBinder extends
|
||||
return stringWriter.getBuffer().toString();
|
||||
}
|
||||
|
||||
/**
|
||||
* Set a {@link ConsumerConfigCustomizer} for the {@link ConsumerFactory} created inside the binder.
|
||||
* @param consumerConfigCustomizer the consumer config customizer
|
||||
* @deprecated in favor of {@link #addConsumerConfigCustomizer(ConsumerConfigCustomizer)}.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setConsumerConfigCustomizer(ConsumerConfigCustomizer consumerConfigCustomizer) {
|
||||
addConsumerConfigCustomizer(consumerConfigCustomizer);
|
||||
}
|
||||
|
||||
public void addConsumerConfigCustomizer(ConsumerConfigCustomizer consumerConfigCustomizer) {
|
||||
if (consumerConfigCustomizer != null) {
|
||||
this.consumerConfigCustomizers.add(consumerConfigCustomizer);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Set a {@link ProducerConfigCustomizer} for the {@link ProducerFactory} created inside the binder.
|
||||
* @param producerConfigCustomizer the producer config customizer
|
||||
* @deprecated in favor of {@link #addProducerConfigCustomizer(ProducerConfigCustomizer)}.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setProducerConfigCustomizer(ProducerConfigCustomizer producerConfigCustomizer) {
|
||||
addProducerConfigCustomizer(producerConfigCustomizer);
|
||||
}
|
||||
|
||||
public void addProducerConfigCustomizer(ProducerConfigCustomizer producerConfigCustomizer) {
|
||||
if (producerConfigCustomizer != null) {
|
||||
this.producerConfigCustomizers.add(producerConfigCustomizer);
|
||||
|
||||
@@ -1515,7 +1515,6 @@ class KafkaBinderTests extends
|
||||
consumerProperties.setBackOffMaxInterval(150);
|
||||
//When auto commit is disabled, then the record is committed after publishing to DLQ using the manual acknowledgement.
|
||||
// (if DLQ is enabled, which is, in this case).
|
||||
consumerProperties.getExtension().setAutoCommitOffset(false);
|
||||
consumerProperties.getExtension().setEnableDlq(true);
|
||||
|
||||
DirectChannel moduleInputChannel = createBindableChannel("input",
|
||||
|
||||
@@ -136,9 +136,6 @@ class KafkaBinderExtendedPropertiesTest {
|
||||
customKafkaConsumerProperties.getConfiguration().get("value.serializer"))
|
||||
.isEqualTo("BarSerializer.class");
|
||||
|
||||
assertThat(kafkaConsumerProperties.isAckEachRecord()).isEqualTo(true);
|
||||
assertThat(customKafkaConsumerProperties.isAckEachRecord()).isEqualTo(false);
|
||||
|
||||
RebalanceListener rebalanceListener = context.getBean(RebalanceListener.class);
|
||||
assertThat(rebalanceListener.latch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
assertThat(rebalanceListener.bindings.keySet()).contains("standard-in",
|
||||
|
||||
@@ -54,21 +54,6 @@ public class RabbitBinderConfigurationProperties {
|
||||
this.adminAddresses = adminAddresses;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param adminAddresses A comma-separated list of RabbitMQ management plugin URLs.
|
||||
* @deprecated in favor of {@link #setAdminAddresses(String[])}. Will be removed in a
|
||||
* future release.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setAdminAdresses(String[] adminAddresses) {
|
||||
setAdminAddresses(adminAddresses);
|
||||
}
|
||||
|
||||
@Deprecated
|
||||
public String[] getAdminAdresses() {
|
||||
return this.adminAddresses;
|
||||
}
|
||||
|
||||
public String[] getNodes() {
|
||||
return nodes;
|
||||
}
|
||||
|
||||
@@ -180,42 +180,6 @@ public class RabbitConsumerProperties extends RabbitCommonProperties {
|
||||
this.prefetch = prefetch;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return the header patterns.
|
||||
* @deprecated - use {@link #getHeaderPatterns()}.
|
||||
*/
|
||||
@Deprecated
|
||||
public String[] getRequestHeaderPatterns() {
|
||||
return this.headerPatterns;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param requestHeaderPatterns request header patterns
|
||||
* @deprecated - use {@link #setHeaderPatterns(String[])}.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setRequestHeaderPatterns(String[] requestHeaderPatterns) {
|
||||
this.headerPatterns = requestHeaderPatterns;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return the tx size.
|
||||
* @deprecated in favor of {@link #getBatchSize()}
|
||||
*/
|
||||
@Deprecated
|
||||
@Min(value = 1, message = "Tx Size should be greater than zero.")
|
||||
public int getTxSize() {
|
||||
return getBatchSize();
|
||||
}
|
||||
|
||||
/**
|
||||
* @param txSize the tx size
|
||||
* deprecated in favor of {@link #setBatchSize(int)}.
|
||||
*/
|
||||
public void setTxSize(int txSize) {
|
||||
setBatchSize(txSize);
|
||||
}
|
||||
|
||||
@Min(value = 1, message = "Batch Size should be greater than zero.")
|
||||
public int getBatchSize() {
|
||||
return batchSize;
|
||||
|
||||
@@ -163,24 +163,6 @@ public class RabbitProducerProperties extends RabbitCommonProperties {
|
||||
*/
|
||||
private boolean superStream;
|
||||
|
||||
/**
|
||||
* @param requestHeaderPatterns the patterns.
|
||||
* @deprecated - use {@link #setHeaderPatterns(String[])}.
|
||||
*/
|
||||
@Deprecated
|
||||
public void setRequestHeaderPatterns(String[] requestHeaderPatterns) {
|
||||
this.headerPatterns = requestHeaderPatterns;
|
||||
}
|
||||
|
||||
/**
|
||||
* @return the header patterns.
|
||||
* @deprecated - use {@link #getHeaderPatterns()}.
|
||||
*/
|
||||
@Deprecated
|
||||
public String[] getRequestHeaderPatterns() {
|
||||
return this.headerPatterns;
|
||||
}
|
||||
|
||||
public void setCompress(boolean compress) {
|
||||
this.compress = compress;
|
||||
}
|
||||
|
||||
@@ -499,7 +499,7 @@ class RabbitBinderTests extends
|
||||
properties.getExtension().setPrefix("foo.");
|
||||
properties.getExtension().setPrefetch(20);
|
||||
properties.getExtension().setHeaderPatterns(new String[] { "foo" });
|
||||
properties.getExtension().setTxSize(10);
|
||||
properties.getExtension().setBatchSize(10);
|
||||
QuorumConfig quorum = properties.getExtension().getQuorum();
|
||||
quorum.setEnabled(true);
|
||||
quorum.setDeliveryLimit(10);
|
||||
|
||||
Reference in New Issue
Block a user