From 1edbd4abba8612bc3bf3de73a32c5f0d601fb860 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 13 Dec 2022 17:28:34 -0500 Subject: [PATCH] Docs cleanup Remove all references to EnableBinding/StreamListener from the docs Resolves https://github.com/spring-cloud/spring-cloud-stream/issues/2589 --- docs/src/main/asciidoc/kafka/kafka_dlq.adoc | 66 ++++++++----------- .../main/asciidoc/kafka/kafka_overview.adoc | 33 ++++++---- .../main/asciidoc/kafka/kafka_partitions.adoc | 35 +++++----- docs/src/main/asciidoc/preface.adoc | 3 - .../main/asciidoc/rabbit/rabbit_overview.adoc | 22 ++++--- .../asciidoc/rabbit/rabbit_partitions.adoc | 31 +++++---- 6 files changed, 97 insertions(+), 93 deletions(-) diff --git a/docs/src/main/asciidoc/kafka/kafka_dlq.adoc b/docs/src/main/asciidoc/kafka/kafka_dlq.adoc index bccd7b404..efdc1de32 100644 --- a/docs/src/main/asciidoc/kafka/kafka_dlq.adoc +++ b/docs/src/main/asciidoc/kafka/kafka_dlq.adoc @@ -87,10 +87,9 @@ spring.cloud.stream.kafka.binder.headers=x-retries ---- .Application -[source, java] +[source,java] ---- @SpringBootApplication -@EnableBinding(TwoOutputProcessor.class) public class ReRouteDlqKApplication implements CommandLineRunner { private static final String X_RETRIES_HEADER = "x-retries"; @@ -102,37 +101,38 @@ public class ReRouteDlqKApplication implements CommandLineRunner { private final AtomicInteger processed = new AtomicInteger(); @Autowired - private MessageChannel parkingLot; + private StreamBridge streamBridge; - @StreamListener(Processor.INPUT) - @SendTo(Processor.OUTPUT) - public Message reRoute(Message failed) { - processed.incrementAndGet(); - Integer retries = failed.getHeaders().get(X_RETRIES_HEADER, Integer.class); - if (retries == null) { - System.out.println("First retry for " + failed); - return MessageBuilder.fromMessage(failed) - .setHeader(X_RETRIES_HEADER, new Integer(1)) - .setHeader(BinderHeaders.PARTITION_OVERRIDE, - failed.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)) - .build(); - } - else if (retries.intValue() < 3) { - System.out.println("Another retry for " + failed); - return MessageBuilder.fromMessage(failed) - .setHeader(X_RETRIES_HEADER, new Integer(retries.intValue() + 1)) - .setHeader(BinderHeaders.PARTITION_OVERRIDE, - failed.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)) - .build(); - } - else { - System.out.println("Retries exhausted for " + failed); - parkingLot.send(MessageBuilder.fromMessage(failed) + @Bean + public Function, Message> reRoute() { + return failed -> { + processed.incrementAndGet(); + Integer retries = failed.getHeaders().get(X_RETRIES_HEADER, Integer.class); + if (retries == null) { + System.out.println("First retry for " + failed); + return MessageBuilder.fromMessage(failed) + .setHeader(X_RETRIES_HEADER, 1) + .setHeader(BinderHeaders.PARTITION_OVERRIDE, + failed.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)) + .build(); + } + else if (retries < 3) { + System.out.println("Another retry for " + failed); + return MessageBuilder.fromMessage(failed) + .setHeader(X_RETRIES_HEADER, retries + 1) + .setHeader(BinderHeaders.PARTITION_OVERRIDE, + failed.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)) + .build(); + } + else { + System.out.println("Retries exhausted for " + failed); + streamBridge.send("parkingLot", MessageBuilder.fromMessage(failed) .setHeader(BinderHeaders.PARTITION_OVERRIDE, failed.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID)) .build()); - } - return null; + } + return null; + }; } @Override @@ -146,13 +146,5 @@ public class ReRouteDlqKApplication implements CommandLineRunner { } } } - - public interface TwoOutputProcessor extends Processor { - - @Output("parkingLot") - MessageChannel parkingLot(); - - } - } ---- diff --git a/docs/src/main/asciidoc/kafka/kafka_overview.adoc b/docs/src/main/asciidoc/kafka/kafka_overview.adoc index fa20c9f9c..46eacfdb8 100644 --- a/docs/src/main/asciidoc/kafka/kafka_overview.adoc +++ b/docs/src/main/asciidoc/kafka/kafka_overview.adoc @@ -527,21 +527,21 @@ Use the corresponding input channel name for your example. [source] ---- @SpringBootApplication -@EnableBinding(Sink.class) public class ManuallyAcknowdledgingConsumer { public static void main(String[] args) { SpringApplication.run(ManuallyAcknowdledgingConsumer.class, args); } - @StreamListener(Sink.INPUT) - public void process(Message message) { - Acknowledgment acknowledgment = message.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); - if (acknowledgment != null) { + @Bean + public Consumer> process() { + return message -> { + Acknowledgment acknowledgment = message.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); + if (acknowledgment != null) { System.out.println("Acknowledgment provided"); acknowledgment.acknowledge(); - } - } + } + }; } ---- @@ -783,16 +783,23 @@ More details on how to suppress meters selectively can be found https://micromet === Tombstone Records (null record values) When using compacted topics, a record with a `null` value (also called a tombstone record) represents the deletion of a key. -To receive such messages in a `@StreamListener` method, the parameter must be marked as not required to receive a `null` value argument. +To receive such messages in a Spring Cloud Stream function, you can use the following strategy. ==== [source, java] ---- -@StreamListener(Sink.INPUT) -public void in(@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) byte[] key, - @Payload(required = false) Customer customer) { - // customer is null if a tombstone record - ... +@Bean +public Function, String> myFunction() { + return value -> { + Object v = value.getPayload(); + String className = v.getClass().getName(); + if (className.isEqualTo("org.springframework.kafka.support.KafkaNull")) { + // this is a tombstone record + } + else { + // continue with processing + } + }; } ---- ==== diff --git a/docs/src/main/asciidoc/kafka/kafka_partitions.adoc b/docs/src/main/asciidoc/kafka/kafka_partitions.adoc index 178e062da..94cde0028 100644 --- a/docs/src/main/asciidoc/kafka/kafka_partitions.adoc +++ b/docs/src/main/asciidoc/kafka/kafka_partitions.adoc @@ -6,10 +6,9 @@ Sometimes it is advantageous to send data to specific partitions -- for example, The following example shows how to configure the producer and consumer side: -[source, java] +[source,java] ---- @SpringBootApplication -@EnableBinding(Source.class) public class KafkaPartitionProducerApplication { private static final Random RANDOM = new Random(System.currentTimeMillis()); @@ -27,13 +26,15 @@ public class KafkaPartitionProducerApplication { .run(args); } - @InboundChannelAdapter(channel = Source.OUTPUT, poller = @Poller(fixedRate = "5000")) - public Message generate() { - String value = data[RANDOM.nextInt(data.length)]; - System.out.println("Sending: " + value); - return MessageBuilder.withPayload(value) - .setHeader("partitionKey", value) - .build(); + @Bean + public Supplier> generate() { + return () -> { + String value = data[RANDOM.nextInt(data.length)]; + System.out.println("Sending: " + value); + return MessageBuilder.withPayload(value) + .setHeader("partitionKey", value) + .build(); + }; } } @@ -46,7 +47,7 @@ spring: cloud: stream: bindings: - output: + generate-out-0: destination: partitioned.topic producer: partition-key-expression: headers['partitionKey'] @@ -66,10 +67,9 @@ Kafka allocates partitions across the instances. The following Spring Boot application listens to a Kafka stream and prints (to the console) the partition ID to which each message goes: -[source, java] +[source,java] ---- @SpringBootApplication -@EnableBinding(Sink.class) public class KafkaPartitionConsumerApplication { public static void main(String[] args) { @@ -78,9 +78,12 @@ public class KafkaPartitionConsumerApplication { .run(args); } - @StreamListener(Sink.INPUT) - public void listen(@Payload String in, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int partition) { - System.out.println(in + " received from partition " + partition); + @Bean + public Consumer> listen() { + return message -> { + int partition =- message.getHeaders().get(KafkaHeaders.RECEIVED_PARTITION_ID); + System.out.println(in + " received from partition " + partition); + }; } } @@ -93,7 +96,7 @@ spring: cloud: stream: bindings: - input: + listen-in-0: destination: partitioned.topic group: myGroup ---- diff --git a/docs/src/main/asciidoc/preface.adoc b/docs/src/main/asciidoc/preface.adoc index f5000842e..57d62102d 100644 --- a/docs/src/main/asciidoc/preface.adoc +++ b/docs/src/main/asciidoc/preface.adoc @@ -165,9 +165,6 @@ For backward compatibility you can still bring `spring-cloud-stream-reactive` from previous versions. - _Test support binder_ `spring-cloud-stream-test-support` with MessageCollector in favor of a new test binder. See <> for more details. - _@StreamMessageConverter_ - deprecated as it is no longer required. -- The `original-content-type` header references have been removed after it's been deprecated in v2.0. -This is primarily for function-based programming model. For StreamListener it would still be required and thus will stay until we deprecate and eventually discontinue StreamListener -and annotation-based programming model. [[spel-and-streaming-data]] diff --git a/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc b/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc index a7c8b4e6a..71b4f2852 100644 --- a/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc +++ b/docs/src/main/asciidoc/rabbit/rabbit_overview.adoc @@ -1271,23 +1271,25 @@ When messages are rejected, they are routed to the DLQ. After 5 seconds, the message expires and is routed to the original queue by using the queue name as the routing key, as shown in the following example: .Spring Boot application -[source, java] +[source,java] ---- @SpringBootApplication -@EnableBinding(Sink.class) public class XDeathApplication { public static void main(String[] args) { SpringApplication.run(XDeathApplication.class, args); } - @StreamListener(Sink.INPUT) - public void listen(String in, @Header(name = "x-death", required = false) Map death) { - if (death != null && death.get("count").equals(3L)) { - // giving up - don't send to DLX - throw new ImmediateAcknowledgeAmqpException("Failed after 4 attempts"); - } - throw new AmqpRejectAndDontRequeueException("failed"); + @Bean + public Consumer> listen() { + return message -> { + Map death = message.getHeaders().get("x-death"); + if (death != null && death.get("count").equals(3L)) { + // giving up - don't send to DLX + throw new ImmediateAcknowledgeAmqpException("Failed after 4 attempts"); + } + throw new AmqpRejectAndDontRequeueException("failed"); + }; } } @@ -1346,7 +1348,7 @@ The health indicator for Rabbit binder delegates to the one provided from Spring For more information on this, see https://docs.spring.io/spring-boot/docs/current/reference/htmlsingle/#actuator.endpoints.health.auto-configured-health-indicators[this]. You can disable this health indicator at the binder level by using the property - `management.health.binders.enabled` and set this to `false`. -In the case of multibinder environements, this has to be set on the binder's environment properties. +In the case of multi-binder environments, this has to be set on the binder's environment properties. When the health indicator is disabled, you should see something like the below in the health actuator endpoint: diff --git a/docs/src/main/asciidoc/rabbit/rabbit_partitions.adoc b/docs/src/main/asciidoc/rabbit/rabbit_partitions.adoc index 0be97700f..e9c64e8a4 100644 --- a/docs/src/main/asciidoc/rabbit/rabbit_partitions.adoc +++ b/docs/src/main/asciidoc/rabbit/rabbit_partitions.adoc @@ -12,7 +12,6 @@ The following Java and YAML examples show how to configure the producer: [source, java] ---- @SpringBootApplication -@EnableBinding(Source.class) public class RabbitPartitionProducerApplication { private static final Random RANDOM = new Random(System.currentTimeMillis()); @@ -30,13 +29,15 @@ public class RabbitPartitionProducerApplication { .run(args); } - @InboundChannelAdapter(channel = Source.OUTPUT, poller = @Poller(fixedRate = "5000")) - public Message generate() { - String value = data[RANDOM.nextInt(data.length)]; - System.out.println("Sending: " + value); - return MessageBuilder.withPayload(value) - .setHeader("partitionKey", value) - .build(); + @Bean + public Supplier> generate() { + return () -> { + String value = data[RANDOM.nextInt(data.length)]; + System.out.println("Sending: " + value); + return MessageBuilder.withPayload(value) + .setHeader("partitionKey", value) + .build(); + }; } } @@ -49,7 +50,7 @@ public class RabbitPartitionProducerApplication { cloud: stream: bindings: - output: + generate-out-0: destination: partitioned.destination producer: partitioned: true @@ -87,7 +88,6 @@ The following Java and YAML examples continue the previous examples and show how [source, java] ---- @SpringBootApplication -@EnableBinding(Sink.class) public class RabbitPartitionConsumerApplication { public static void main(String[] args) { @@ -96,9 +96,12 @@ public class RabbitPartitionConsumerApplication { .run(args); } - @StreamListener(Sink.INPUT) - public void listen(@Payload String in, @Header(AmqpHeaders.CONSUMER_QUEUE) String queue) { - System.out.println(in + " received from queue " + queue); + @Bean + public Consumer> listen() { + return message -> { + String queue =- message.getHeaders().get(AmqpHeaders.CONSUMER_QUEUE); + System.out.println(in + " received from queue " + queue); + }; } } @@ -111,7 +114,7 @@ public class RabbitPartitionConsumerApplication { cloud: stream: bindings: - input: + listen-in-0: destination: partitioned.destination group: myGroup consumer: