Polish "Error Handling Docs"
This commit is contained in:
@@ -860,7 +860,7 @@ If you want to specify some advanced backoff options for ack timeout with differ
|
||||
----
|
||||
@EnablePulsar
|
||||
@Configuration
|
||||
static class AckTimeoutRedeliveryConfig {
|
||||
class AckTimeoutRedeliveryConfig {
|
||||
|
||||
@PulsarListener(subscriptionName = "withAckTimeoutRedeliveryBackoffSubscription",
|
||||
topics = "withAckTimeoutRedeliveryBackoff-test-topic",
|
||||
@@ -871,7 +871,7 @@ static class AckTimeoutRedeliveryConfig {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public RedeliveryBackoff ackTimeoutRedeliveryBackoff() {
|
||||
RedeliveryBackoff ackTimeoutRedeliveryBackoff() {
|
||||
return MultiplierRedeliveryBackoff.builder().minDelayMs(1000).maxDelayMs(10 * 1000).multiplier(2)
|
||||
.build();
|
||||
}
|
||||
@@ -909,7 +909,7 @@ Here is an example:
|
||||
----
|
||||
@EnablePulsar
|
||||
@Configuration
|
||||
static class NegativeAckRedeliveryConfig {
|
||||
class NegativeAckRedeliveryConfig {
|
||||
|
||||
@PulsarListener(subscriptionName = "withNegRedeliveryBackoffSubscription",
|
||||
topics = "withNegRedeliveryBackoff-test-topic", negativeAckRedeliveryBackoff = "redeliveryBackoff",
|
||||
@@ -919,7 +919,7 @@ static class NegativeAckRedeliveryConfig {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public RedeliveryBackoff redeliveryBackoff() {
|
||||
RedeliveryBackoff redeliveryBackoff() {
|
||||
return MultiplierRedeliveryBackoff.builder().minDelayMs(1000).maxDelayMs(10 * 1000).multiplier(2)
|
||||
.build();
|
||||
}
|
||||
@@ -940,7 +940,7 @@ Let us see some details around this feature in action by inspecting some code sn
|
||||
----
|
||||
@EnablePulsar
|
||||
@Configuration
|
||||
static class DeadLetterPolicyConfig {
|
||||
class DeadLetterPolicyConfig {
|
||||
|
||||
@PulsarListener(id = "deadLetterPolicyListener", subscriptionName = "deadLetterPolicySubscription",
|
||||
topics = "topic-with-dlp", deadLetterPolicy = "deadLetterPolicy",
|
||||
@@ -997,10 +997,10 @@ Let us see some details by examining a few code snippets.
|
||||
----
|
||||
@EnablePulsar
|
||||
@Configuration
|
||||
static class PulsarConsumerErrorHandlerConfig {
|
||||
class PulsarConsumerErrorHandlerConfig {
|
||||
|
||||
@Bean
|
||||
public PulsarConsumerErrorHandler<String> pulsarConsumerErrorHandler(
|
||||
@Bean
|
||||
PulsarConsumerErrorHandler<String> pulsarConsumerErrorHandler(
|
||||
PulsarTemplate<String> pulsarTemplate) {
|
||||
return new DefaultPulsarConsumerErrorHandler<>(
|
||||
new PulsarDeadLetterPublishingRecoverer<>(pulsarTemplate, (c, m) -> "my-foo-dlt"), new FixedBackOff(100, 10));
|
||||
@@ -1083,18 +1083,18 @@ Towards this extent, you can use a `PulsarConsumerErrorHandler` with the followi
|
||||
[source, java]
|
||||
----
|
||||
@Bean
|
||||
public PulsarConsumerErrorHandler<Integer> pulsarConsumerErrorHandler(PulsarClient pulsarClient) {
|
||||
PulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, Map.of());
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
PulsarConsumerErrorHandler<Integer> pulsarConsumerErrorHandler(PulsarClient pulsarClient) {
|
||||
PulsarProducerFactory<Integer> pulsarProducerFactory = new DefaultPulsarProducerFactory<>(pulsarClient, Map.of());
|
||||
PulsarTemplate<Integer> pulsarTemplate = new PulsarTemplate<>(pulsarProducerFactory);
|
||||
|
||||
BiFunction<Consumer<?>, Message<?>, String> destinationResolver =
|
||||
(c, m) -> "my-foo-dlt";
|
||||
BiFunction<Consumer<?>, Message<?>, String> destinationResolver =
|
||||
(c, m) -> "my-foo-dlt";
|
||||
|
||||
final PulsarDeadLetterPublishingRecoverer<Integer> pulsarDeadLetterPublishingRecoverer =
|
||||
new PulsarDeadLetterPublishingRecoverer<>(pulsarTemplate, destinationResolver);
|
||||
PulsarDeadLetterPublishingRecoverer<Integer> pulsarDeadLetterPublishingRecoverer =
|
||||
new PulsarDeadLetterPublishingRecoverer<>(pulsarTemplate, destinationResolver);
|
||||
|
||||
return new DefaultPulsarConsumerErrorHandler<>(pulsarDeadLetterPublishingRecoverer,
|
||||
new FixedBackOff(100, 5));
|
||||
return new DefaultPulsarConsumerErrorHandler<>(pulsarDeadLetterPublishingRecoverer,
|
||||
new FixedBackOff(100, 5));
|
||||
}
|
||||
----
|
||||
====
|
||||
@@ -1121,17 +1121,17 @@ First, let us look at a batch `PulsarListener` method.
|
||||
@PulsarListener(subscriptionName = "batch-demo-5-sub", topics = "batch-demo-4", batch = true, concurrency = "3",
|
||||
subscriptionType = SubscriptionType.Failover,
|
||||
pulsarConsumerErrorHandler = "pulsarConsumerErrorHandler", ackMode = AckMode.MANUAL)
|
||||
public void listen(List<Message<Integer>> data, Consumer<Integer> consumer, Acknowledgment acknowledgment) {
|
||||
void listen(List<Message<Integer>> data, Consumer<Integer> consumer, Acknowledgment acknowledgment) {
|
||||
for (Message<Integer> datum : data) {
|
||||
if (datum.getValue() == 5) {
|
||||
if (datum.getValue() == 5) {
|
||||
throw new PulsarBatchListenerFailedException("failed", datum);
|
||||
}
|
||||
acknowledgement.acknowledge(datum.getMessageId());
|
||||
}
|
||||
}
|
||||
acknowledgement.acknowledge(datum.getMessageId());
|
||||
}
|
||||
}
|
||||
|
||||
@Bean
|
||||
public PulsarConsumerErrorHandler<String> pulsarConsumerErrorHandler(
|
||||
PulsarConsumerErrorHandler<String> pulsarConsumerErrorHandler(
|
||||
PulsarTemplate<String> pulsarTemplate) {
|
||||
return new DefaultPulsarConsumerErrorHandler<>(
|
||||
new PulsarDeadLetterPublishingRecoverer<>(pulsarTemplate, (c, m) -> "my-foo-dlt"), new FixedBackOff(100, 10));
|
||||
|
||||
Reference in New Issue
Block a user