From a3799e5f7a9522efca00326f7c2bc26830e4f114 Mon Sep 17 00:00:00 2001 From: dibulidohu Date: Wed, 4 Aug 2021 11:50:34 +0800 Subject: [PATCH] fix retryLetterProducer usr auto_chema --- .../java/org/apache/pulsar/client/impl/ConsumerImpl.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 17d804390ec03..879f9367c434f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -170,7 +170,7 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private volatile CompletableFuture> deadLetterProducer; - private volatile Producer retryLetterProducer; + private volatile Producer retryLetterProducer; private final ReadWriteLock createProducerLock = new ReentrantReadWriteLock(); protected volatile boolean paused; @@ -559,7 +559,7 @@ protected CompletableFuture doReconsumeLater(Message message, AckType a createProducerLock.writeLock().lock(); try { if (retryLetterProducer == null) { - retryLetterProducer = client.newProducer(schema) + retryLetterProducer = client.newProducer(Schema.AUTO_PRODUCE_BYTES(schema)) .topic(this.deadLetterPolicy.getRetryLetterTopic()) .enableBatching(false) .blockIfQueueFull(false) @@ -612,8 +612,9 @@ protected CompletableFuture doReconsumeLater(Message message, AckType a return null; }); } else { - TypedMessageBuilder typedMessageBuilderNew = retryLetterProducer.newMessage() - .value(retryMessage.getValue()) + assert retryMessage != null; + TypedMessageBuilder typedMessageBuilderNew = retryLetterProducer.newMessage() + .value(retryMessage.getData()) .properties(propertiesMap); if (delayTime > 0) { typedMessageBuilderNew.deliverAfter(delayTime, unit);