diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 311e82dde..e94f76a21 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -265,6 +265,10 @@ The replication factor to use when provisioning topics. Overrides the binder-wid Ignored if `replicas-assignments` is present. + Default: none (the binder-wide default of 1 is used). +pollTimeout:: +Timeout used for polling in pollable consumers. ++ +Default: 5 seconds. ==== Consuming Batches diff --git a/pom.xml b/pom.xml index f61e268dd..a776c9bc9 100644 --- a/pom.xml +++ b/pom.xml @@ -13,8 +13,9 @@ 1.8 2.3.0.BUILD-SNAPSHOT - 3.2.0.M4 + 3.2.0.BUILD-SNAPSHOT 2.3.0 + 1.0.0.BUILD-SNAPSHOT 3.0.0.BUILD-SNAPSHOT true true @@ -113,8 +114,8 @@ org.springframework.cloud - spring-cloud-stream-schema - ${spring-cloud-stream.version} + spring-cloud-schema-registry-client + ${spring-cloud-schema-registry.version} test diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 8d9409a0c..5641fbb14 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2018 the original author or authors. + * Copyright 2016-2019 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -120,6 +120,11 @@ public class KafkaConsumerProperties { private KafkaTopicProperties topic = new KafkaTopicProperties(); + /** + * Timeout used for polling in pollable consumers. + */ + private long pollTimeout = org.springframework.kafka.listener.ConsumerProperties.DEFAULT_POLL_TIMEOUT; + public boolean isAckEachRecord() { return this.ackEachRecord; } @@ -291,4 +296,11 @@ public class KafkaConsumerProperties { this.topic = topic; } + public long getPollTimeout() { + return this.pollTimeout; + } + + public void setPollTimeout(long pollTimeout) { + this.pollTimeout = pollTimeout; + } } diff --git a/spring-cloud-stream-binder-kafka-streams/pom.xml b/spring-cloud-stream-binder-kafka-streams/pom.xml index 4e8a2f9f8..d474fd99c 100644 --- a/spring-cloud-stream-binder-kafka-streams/pom.xml +++ b/spring-cloud-stream-binder-kafka-streams/pom.xml @@ -88,7 +88,7 @@ org.springframework.cloud - spring-cloud-stream-schema + spring-cloud-schema-registry-client test diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java index 1b866e1df..00c745a3b 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/PerRecordAvroContentTypeTests.java @@ -37,13 +37,13 @@ import org.junit.Test; import org.springframework.boot.SpringApplication; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.Input; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.annotation.StreamMessageConverter; import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.cloud.stream.binder.kafka.streams.integration.utils.TestAvroSerializer; -import org.springframework.cloud.stream.schema.avro.AvroSchemaMessageConverter; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; @@ -60,6 +60,7 @@ import org.springframework.util.MimeTypeUtils; import static org.assertj.core.api.Assertions.assertThat; + /** * @author Soby Chacko */ @@ -174,9 +175,8 @@ public class PerRecordAvroContentTypeTests { } @Bean - @StreamMessageConverter public MessageConverter sensorMessageConverter() throws IOException { - return new AvroSchemaMessageConverter(); + return new AvroSchemaMessageConverter(new AvroSchemaServiceManagerImpl()); } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java index b771ed698..6bbf3180e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/integration/utils/TestAvroSerializer.java @@ -21,7 +21,8 @@ import java.util.Map; import org.apache.kafka.common.serialization.Serializer; -import org.springframework.cloud.stream.schema.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.support.MessageBuilder; @@ -44,7 +45,7 @@ public class TestAvroSerializer implements Serializer { @Override public byte[] serialize(String topic, S data) { - AvroSchemaMessageConverter avroSchemaMessageConverter = new AvroSchemaMessageConverter(); + AvroSchemaMessageConverter avroSchemaMessageConverter = new AvroSchemaMessageConverter(new AvroSchemaServiceManagerImpl()); Message message = MessageBuilder.withPayload(data).build(); Map headers = new HashMap<>(message.getHeaders()); headers.put(MessageHeaders.CONTENT_TYPE, "application/avro"); diff --git a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java index 9ba84b187..2f3e0345a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java +++ b/spring-cloud-stream-binder-kafka-streams/src/test/java/org/springframework/cloud/stream/binder/kafka/streams/serde/CompositeNonNativeSerdeTest.java @@ -27,8 +27,9 @@ import com.example.Sensor; import com.fasterxml.jackson.databind.ObjectMapper; import org.junit.Test; +import org.springframework.cloud.schema.registry.avro.AvroSchemaMessageConverter; +import org.springframework.cloud.schema.registry.avro.AvroSchemaServiceManagerImpl; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; -import org.springframework.cloud.stream.schema.avro.AvroSchemaMessageConverter; import org.springframework.messaging.converter.MessageConverter; import static org.assertj.core.api.Assertions.assertThat; @@ -51,7 +52,7 @@ public class CompositeNonNativeSerdeTest { sensor.setTemperature(random.nextFloat() * 50); List messageConverters = new ArrayList<>(); - messageConverters.add(new AvroSchemaMessageConverter()); + messageConverters.add(new AvroSchemaMessageConverter(new AvroSchemaServiceManagerImpl())); CompositeMessageConverterFactory compositeMessageConverterFactory = new CompositeMessageConverterFactory( messageConverters, new ObjectMapper()); CompositeNonNativeSerde compositeNonNativeSerde = new CompositeNonNativeSerde( diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index fbf2f5f0d..0cfd8899c 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -97,6 +97,7 @@ import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.AbstractMessageListenerContainer; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ConsumerAwareRebalanceListener; +import org.springframework.kafka.listener.ConsumerProperties; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.support.DefaultKafkaHeaderMapper; import org.springframework.kafka.support.KafkaHeaderMapper; @@ -778,53 +779,34 @@ public class KafkaMessageChannelBinder extends @Override protected PolledConsumerResources createPolledConsumerResources(String name, String group, ConsumerDestination destination, - ExtendedConsumerProperties consumerProperties) { + ExtendedConsumerProperties extendedConsumerProperties) { boolean anonymous = !StringUtils.hasText(group); - Assert.isTrue(!anonymous || !consumerProperties.getExtension().isEnableDlq(), + final KafkaConsumerProperties extension = extendedConsumerProperties.getExtension(); + Assert.isTrue(!anonymous || !extension.isEnableDlq(), "DLQ support is not available for anonymous subscriptions"); String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group; final ConsumerFactory consumerFactory = createKafkaConsumerFactory( - anonymous, consumerGroup, consumerProperties); - String[] topics = consumerProperties.isMultiplex() + anonymous, consumerGroup, extendedConsumerProperties); + String[] topics = extendedConsumerProperties.isMultiplex() ? StringUtils.commaDelimitedListToStringArray(destination.getName()) : new String[] { destination.getName() }; for (int i = 0; i < topics.length; i++) { topics[i] = topics[i].trim(); } - KafkaMessageSource source = new KafkaMessageSource<>(consumerFactory, - topics); - source.setMessageConverter(getMessageConverter(consumerProperties)); - source.setRawMessageHeader(consumerProperties.getExtension().isEnableDlq()); + final ConsumerProperties consumerProperties = new ConsumerProperties(topics); + String clientId = name; - if (consumerProperties.getExtension().getConfiguration() + if (extension.getConfiguration() .containsKey(ConsumerConfig.CLIENT_ID_CONFIG)) { - clientId = consumerProperties.getExtension().getConfiguration() + clientId = extension.getConfiguration() .get(ConsumerConfig.CLIENT_ID_CONFIG); } - source.setClientId(clientId); - if (!consumerProperties.isMultiplex()) { - // I copied this from the regular consumer - it looks bogus to me - includes - // all partitions - // not just the ones this binding is listening to; doesn't seem right for a - // health check. - Collection partitionInfos = getPartitionInfo( - destination.getName(), consumerProperties, consumerFactory, -1); - this.topicsInUse.put(destination.getName(), - new TopicInformation(consumerGroup, partitionInfos, false)); - } - else { - for (int i = 0; i < topics.length; i++) { - Collection partitionInfos = getPartitionInfo(topics[i], - consumerProperties, consumerFactory, -1); - this.topicsInUse.put(topics[i], - new TopicInformation(consumerGroup, partitionInfos, false)); - } - } + consumerProperties.setClientId(clientId); - source.setRebalanceListener(new ConsumerRebalanceListener() { + consumerProperties.setConsumerRebalanceListener(new ConsumerRebalanceListener() { @Override public void onPartitionsRevoked(Collection partitions) { @@ -837,9 +819,36 @@ public class KafkaMessageChannelBinder extends } }); + + consumerProperties.setPollTimeout(extension.getPollTimeout()); + + KafkaMessageSource source = new KafkaMessageSource<>(consumerFactory, + consumerProperties); + source.setMessageConverter(getMessageConverter(extendedConsumerProperties)); + source.setRawMessageHeader(extension.isEnableDlq()); + + if (!extendedConsumerProperties.isMultiplex()) { + // I copied this from the regular consumer - it looks bogus to me - includes + // all partitions + // not just the ones this binding is listening to; doesn't seem right for a + // health check. + Collection partitionInfos = getPartitionInfo( + destination.getName(), extendedConsumerProperties, consumerFactory, -1); + this.topicsInUse.put(destination.getName(), + new TopicInformation(consumerGroup, partitionInfos, false)); + } + else { + for (int i = 0; i < topics.length; i++) { + Collection partitionInfos = getPartitionInfo(topics[i], + extendedConsumerProperties, consumerFactory, -1); + this.topicsInUse.put(topics[i], + new TopicInformation(consumerGroup, partitionInfos, false)); + } + } + getMessageSourceCustomizer().configure(source, destination.getName(), group); return new PolledConsumerResources(source, registerErrorInfrastructure( - destination, group, consumerProperties, true)); + destination, group, extendedConsumerProperties, true)); } @Override diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index ff4a40b11..7fd451cb4 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -2942,6 +2942,7 @@ public class KafkaBinderTests extends this.messageConverter); ExtendedConsumerProperties consumerProps = createConsumerProperties(); consumerProps.setMultiplex(true); + consumerProps.getExtension().setPollTimeout(1); Binding> binding = binder.bindPollableConsumer( "pollable,anotherOne", "group-polledConsumer", inboundBindTarget, consumerProps); @@ -3028,6 +3029,7 @@ public class KafkaBinderTests extends PollableSource inboundBindTarget = new DefaultPollableMessageSource( this.messageConverter); ExtendedConsumerProperties properties = createConsumerProperties(); + properties.getExtension().setPollTimeout(1); properties.setMaxAttempts(2); properties.setBackOffInitialInterval(0); properties.getExtension().setEnableDlq(true);