Polish "Batch receive properties not copied properly"

This commit is contained in:
Chris Bono
2022-09-16 20:24:50 -05:00
committed by Chris Bono
parent 9d48b8df1b
commit db3782160e
8 changed files with 32 additions and 51 deletions

View File

@@ -61,7 +61,7 @@ public class PulsarAnnotationDrivenConfiguration {
map.from(properties::getSchemaType).to(containerProperties::setSchemaType);
map.from(properties::getAckMode).to(containerProperties::setAckMode);
map.from(properties::getBatchTimeoutMillis).to(containerProperties::setBatchTimeout);
map.from(properties::getBatchTimeoutMillis).to(containerProperties::setBatchTimeoutMillis);
map.from(properties::getMaxNumBytes).to(containerProperties::setMaxNumBytes);
map.from(properties::getMaxNumMessages).to(containerProperties::setMaxNumMessages);

View File

@@ -210,7 +210,7 @@ class PulsarAutoConfigurationTests {
.getBean(ConcurrentPulsarListenerContainerFactory.class).extracting("containerProperties")
.hasFieldOrPropertyWithValue("maxNumMessages", 10)
.hasFieldOrPropertyWithValue("maxNumBytes", 101)
.hasFieldOrPropertyWithValue("batchTimeout", 50)));
.hasFieldOrPropertyWithValue("batchTimeoutMillis", 50)));
}
@Nested

View File

@@ -16,14 +16,10 @@
package org.springframework.pulsar.config;
import java.util.Arrays;
import org.apache.commons.logging.LogFactory;
import org.apache.pulsar.client.api.Schema;
import org.springframework.beans.BeanWrapper;
import org.springframework.beans.BeansException;
import org.springframework.beans.PropertyAccessorFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.context.ApplicationContext;
import org.springframework.context.ApplicationContextAware;
@@ -136,35 +132,32 @@ public abstract class AbstractPulsarListenerContainerFactory<C extends AbstractP
protected abstract C createContainerInstance(PulsarListenerEndpoint endpoint);
private void configureEndpoint(AbstractPulsarListenerEndpoint<C> aplEndpoint) {
if (aplEndpoint.getBatchListener() == null) {
JavaUtils.INSTANCE.acceptIfNotNull(this.batchListener, aplEndpoint::setBatchListener);
}
}
protected void initializeContainer(C instance, PulsarListenerEndpoint endpoint) {
PulsarContainerProperties properties = instance.getPulsarContainerProperties();
// BeanUtils.copyProperties(this.containerProperties, properties, "topics",
// "messageListener",
// "batchListener", "subscriptionName", "subscriptionType", "schema");
PulsarContainerProperties instanceProperties = instance.getPulsarContainerProperties();
if (properties.getSchemaType() == null) {
if (this.containerProperties.getSchemaType() != null) {
properties.setSchemaType(this.containerProperties.getSchemaType());
}
if (instanceProperties.getSchemaType() == null) {
JavaUtils.INSTANCE.acceptIfNotNull(this.containerProperties.getSchemaType(),
instanceProperties::setSchemaType);
}
if (properties.getSchema() == null) {
properties.setSchema(Schema.BYTES);
if (instanceProperties.getSchema() == null) {
instanceProperties.setSchema(Schema.BYTES);
}
if (properties.getSubscriptionType() == null) {
properties.setSubscriptionType(this.containerProperties.getSubscriptionType());
if (instanceProperties.getSubscriptionType() == null) {
instanceProperties.setSubscriptionType(this.containerProperties.getSubscriptionType());
}
if (endpoint.getAckMode() != AckMode.BATCH) {
properties.setAckMode(endpoint.getAckMode());
instanceProperties.setAckMode(endpoint.getAckMode());
}
else if (this.containerProperties.getAckMode() != AckMode.BATCH) {
properties.setAckMode(this.containerProperties.getAckMode());
instanceProperties.setAckMode(this.containerProperties.getAckMode());
}
Boolean autoStart = endpoint.getAutoStartup();
@@ -175,8 +168,9 @@ public abstract class AbstractPulsarListenerContainerFactory<C extends AbstractP
instance.setAutoStartup(this.autoStartup);
}
copyProperties(this.containerProperties, instance.getContainerProperties(),
Arrays.asList("maxNumMessages", "maxNumBytes", "batchTimeout"));
instanceProperties.setMaxNumMessages(this.containerProperties.getMaxNumMessages());
instanceProperties.setMaxNumBytes(this.containerProperties.getMaxNumBytes());
instanceProperties.setBatchTimeoutMillis(this.containerProperties.getBatchTimeoutMillis());
JavaUtils.INSTANCE.acceptIfNotNull(this.phase, instance::setPhase)
.acceptIfNotNull(this.applicationContext, instance::setApplicationContext)
@@ -185,17 +179,4 @@ public abstract class AbstractPulsarListenerContainerFactory<C extends AbstractP
instance.getContainerProperties()::setPulsarConsumerProperties);
}
/**
* Copy a list of properties from the source object to the target.
* @param source Object source to copy from
* @param target Object target to copy to
* @param requestProperties list of properties to copy from source to target
*/
public static void copyProperties(Object source, Object target, Iterable<String> requestProperties) {
BeanWrapper wrappedSource = PropertyAccessorFactory.forBeanPropertyAccess(source);
BeanWrapper wrappedTarget = PropertyAccessorFactory.forBeanPropertyAccess(target);
requestProperties.forEach(p -> wrappedTarget.setPropertyValue(p, wrappedSource.getPropertyValue(p)));
}
}

View File

@@ -193,7 +193,7 @@ public class DefaultPulsarMessageListenerContainer<T> extends AbstractPulsarMess
final BatchReceivePolicy batchReceivePolicy = new BatchReceivePolicy.Builder()
.maxNumMessages(pulsarContainerProperties.getMaxNumMessages())
.maxNumBytes(pulsarContainerProperties.getMaxNumBytes())
.timeout(pulsarContainerProperties.getBatchTimeout(), TimeUnit.MILLISECONDS).build();
.timeout(pulsarContainerProperties.getBatchTimeoutMillis(), TimeUnit.MILLISECONDS).build();
this.consumer = getPulsarConsumerFactory().createConsumer(
(Schema) pulsarContainerProperties.getSchema(), batchReceivePolicy, propertiesToConsumer);
Assert.state(this.consumer != null, "Unable to create a consumer");

View File

@@ -58,7 +58,7 @@ public class PulsarContainerProperties {
private int maxNumBytes = 10 * 1024 * 1024;
private int batchTimeout = 100;
private int batchTimeoutMillis = 100;
private boolean batchListener;
@@ -116,12 +116,12 @@ public class PulsarContainerProperties {
this.maxNumBytes = maxNumBytes;
}
public int getBatchTimeout() {
return this.batchTimeout;
public int getBatchTimeoutMillis() {
return this.batchTimeoutMillis;
}
public void setBatchTimeout(int batchTimeout) {
this.batchTimeout = batchTimeout;
public void setBatchTimeoutMillis(int batchTimeoutMillis) {
this.batchTimeoutMillis = batchTimeoutMillis;
}
public boolean isBatchListener() {

View File

@@ -271,7 +271,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
CountDownLatch latch = new CountDownLatch(1);
final PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
@@ -319,7 +319,7 @@ class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests {
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchMessageListener<?> pulsarBatchMessageListener = mock(PulsarBatchMessageListener.class);
CountDownLatch latch = new CountDownLatch(1);

View File

@@ -57,7 +57,7 @@ public class ConcurrentPulsarMessageListenerContainerTests {
factory.setPulsarConsumerFactory(pulsarConsumerFactory);
PulsarContainerProperties containerProperties = factory.getContainerProperties();
containerProperties.setBatchTimeout(60_000);
containerProperties.setBatchTimeoutMillis(60_000);
containerProperties.setMaxNumMessages(120);
containerProperties.setMaxNumBytes(32000);
@@ -69,7 +69,7 @@ public class ConcurrentPulsarMessageListenerContainerTests {
final ConcurrentPulsarMessageListenerContainer<String> concurrentContainer = factory
.createListenerContainer(pulsarListenerEndpoint);
final PulsarContainerProperties pulsarContainerProperties = concurrentContainer.getContainerProperties();
assertThat(pulsarContainerProperties.getBatchTimeout()).isEqualTo(60_000);
assertThat(pulsarContainerProperties.getBatchTimeoutMillis()).isEqualTo(60_000);
assertThat(pulsarContainerProperties.getMaxNumMessages()).isEqualTo(120);
assertThat(pulsarContainerProperties.getMaxNumBytes()).isEqualTo(32_000);
}

View File

@@ -223,7 +223,7 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
PulsarBatchAcknowledgingMessageListener.class);
@@ -294,7 +294,7 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
PulsarBatchAcknowledgingMessageListener.class);
@@ -363,7 +363,7 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
PulsarBatchAcknowledgingMessageListener.class);
@@ -432,7 +432,7 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
PulsarBatchAcknowledgingMessageListener.class);
@@ -500,7 +500,7 @@ public class DefaultPulsarConsumerErrorHandlerTests extends AbstractContainerBas
PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties();
pulsarContainerProperties.setMaxNumMessages(10);
pulsarContainerProperties.setBatchTimeout(60_000);
pulsarContainerProperties.setBatchTimeoutMillis(60_000);
pulsarContainerProperties.setBatchListener(true);
final PulsarBatchAcknowledgingMessageListener<?> pulsarBatchMessageListener = mock(
PulsarBatchAcknowledgingMessageListener.class);