|
41 | 41 |
|
42 | 42 | import org.springframework.boot.autoconfigure.pulsar.PulsarProperties.Consumer;
|
43 | 43 | import org.springframework.boot.autoconfigure.pulsar.PulsarProperties.Failover.BackupCluster;
|
44 |
| -import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory; |
45 | 44 | import org.springframework.pulsar.core.PulsarProducerFactory;
|
46 | 45 | import org.springframework.pulsar.core.PulsarTemplate;
|
47 | 46 | import org.springframework.pulsar.listener.PulsarContainerProperties;
|
@@ -263,27 +262,18 @@ void customizeContainerProperties() {
|
263 | 262 | PulsarProperties properties = new PulsarProperties();
|
264 | 263 | properties.getConsumer().getSubscription().setType(SubscriptionType.Shared);
|
265 | 264 | properties.getListener().setSchemaType(SchemaType.AVRO);
|
| 265 | + properties.getListener().setConcurrency(10); |
266 | 266 | properties.getListener().setObservationEnabled(true);
|
267 | 267 | properties.getTransaction().setEnabled(true);
|
268 | 268 | PulsarContainerProperties containerProperties = new PulsarContainerProperties("my-topic-pattern");
|
269 | 269 | new PulsarPropertiesMapper(properties).customizeContainerProperties(containerProperties);
|
270 | 270 | assertThat(containerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Shared);
|
271 | 271 | assertThat(containerProperties.getSchemaType()).isEqualTo(SchemaType.AVRO);
|
| 272 | + assertThat(containerProperties.getConcurrency()).isEqualTo(10); |
272 | 273 | assertThat(containerProperties.isObservationEnabled()).isTrue();
|
273 | 274 | assertThat(containerProperties.transactions().isEnabled()).isTrue();
|
274 | 275 | }
|
275 | 276 |
|
276 |
| - @Test |
277 |
| - void customizeConcurrentPulsarListenerContainerFactory() { |
278 |
| - PulsarProperties properties = new PulsarProperties(); |
279 |
| - properties.getListener().setConcurrency(10); |
280 |
| - ConcurrentPulsarListenerContainerFactory<?> listenerContainerFactory = mock( |
281 |
| - ConcurrentPulsarListenerContainerFactory.class); |
282 |
| - new PulsarPropertiesMapper(properties) |
283 |
| - .customizeConcurrentPulsarListenerContainerFactory(listenerContainerFactory); |
284 |
| - then(listenerContainerFactory).should().setConcurrency(10); |
285 |
| - } |
286 |
| - |
287 | 277 | @Test
|
288 | 278 | @SuppressWarnings("unchecked")
|
289 | 279 | void customizeReaderBuilder() {
|
|
0 commit comments