前往小程序,Get更优阅读体验!
立即前往
首页
学习
活动
专区
工具
TVP
发布
社区首页 >专栏 >聊聊spring kafka的retry

聊聊spring kafka的retry

作者头像
code4it
发布2018-09-17 15:11:47
1.1K0
发布2018-09-17 15:11:47
举报
文章被收录于专栏:码匠的流水账码匠的流水账

本文主要聊一下spring for kafka的retry

AbstractRetryingMessageListenerAdapter

spring-kafka-1.2.3.RELEASE-sources.jar!/org/springframework/kafka/listener/adapter/AbstractRetryingMessageListenerAdapter.java 主要有两个实现类RetryingAcknowledgingMessageListenerAdapter以及RetryingMessageListenerAdapter

RetryingAcknowledgingMessageListenerAdapter

代码语言:javascript
复制
public class RetryingAcknowledgingMessageListenerAdapter<K, V>
        extends AbstractRetryingMessageListenerAdapter<K, V, AcknowledgingMessageListener<K, V>>
        implements AcknowledgingMessageListener<K, V> {

    private final AcknowledgingMessageListener<K, V> delegate;

    /**
     * Construct an instance with the provided template and delegate. The exception will
     * be thrown to the container after retries are exhausted.
     * @param messageListener the listener delegate.
     * @param retryTemplate the template.
     */
    public RetryingAcknowledgingMessageListenerAdapter(AcknowledgingMessageListener<K, V> messageListener,
            RetryTemplate retryTemplate) {
        this(messageListener, retryTemplate, null);
    }

    /**
     * Construct an instance with the provided template, callback and delegate.
     * @param messageListener the listener delegate.
     * @param retryTemplate the template.
     * @param recoveryCallback the recovery callback; if null, the exception will be
     * thrown to the container after retries are exhausted.
     */
    public RetryingAcknowledgingMessageListenerAdapter(AcknowledgingMessageListener<K, V> messageListener,
            RetryTemplate retryTemplate, RecoveryCallback<? extends Object> recoveryCallback) {
        super(messageListener, retryTemplate, recoveryCallback);
        Assert.notNull(messageListener, "'messageListener' cannot be null");
        this.delegate = messageListener;
    }

    @SuppressWarnings("unchecked")
    @Override
    public void onMessage(final ConsumerRecord<K, V> record, final Acknowledgment acknowledgment) {
        getRetryTemplate().execute(new RetryCallback<Object, KafkaException>() {

            @Override
            public Void doWithRetry(RetryContext context) throws KafkaException {
                context.setAttribute(CONTEXT_RECORD, record);
                context.setAttribute(CONTEXT_ACKNOWLEDGMENT, acknowledgment);
                RetryingAcknowledgingMessageListenerAdapter.this.delegate.onMessage(record, acknowledgment);
                return null;
            }

        }, (RecoveryCallback<Object>) getRecoveryCallback());
    }

}

RetryingMessageListenerAdapter

spring-kafka-1.2.3.RELEASE-sources.jar!/org/springframework/kafka/listener/adapter/RetryingMessageListenerAdapter.java

代码语言:javascript
复制
public class RetryingMessageListenerAdapter<K, V>
        extends AbstractRetryingMessageListenerAdapter<K, V, MessageListener<K, V>>
        implements MessageListener<K, V> {

    /**
     * Construct an instance with the provided template and delegate. The exception will
     * be thrown to the container after retries are exhausted.
     * @param messageListener the delegate listener.
     * @param retryTemplate the template.
     */
    public RetryingMessageListenerAdapter(MessageListener<K, V> messageListener, RetryTemplate retryTemplate) {
        this(messageListener, retryTemplate, null);
    }

    /**
     * Construct an instance with the provided template, callback and delegate.
     * @param messageListener the delegate listener.
     * @param retryTemplate the template.
     * @param recoveryCallback the recovery callback; if null, the exception will be
     * thrown to the container after retries are exhausted.
     */
    public RetryingMessageListenerAdapter(MessageListener<K, V> messageListener, RetryTemplate retryTemplate,
            RecoveryCallback<? extends Object> recoveryCallback) {
        super(messageListener, retryTemplate, recoveryCallback);
        Assert.notNull(messageListener, "'messageListener' cannot be null");
    }

    @SuppressWarnings("unchecked")
    @Override
    public void onMessage(final ConsumerRecord<K, V> record) {
        getRetryTemplate().execute(new RetryCallback<Object, KafkaException>() {

            @Override
            public Void doWithRetry(RetryContext context) throws KafkaException {
                context.setAttribute(CONTEXT_RECORD, record);
                RetryingMessageListenerAdapter.this.delegate.onMessage(record);
                return null;
            }

        }, (RecoveryCallback<Object>) getRecoveryCallback());
    }

}

具体是哪种listener呢

spring-kafka-1.2.3.RELEASE-sources.jar!/org/springframework/kafka/config/AbstractKafkaListenerEndpoint.java

代码语言:javascript
复制
private void setupMessageListener(MessageListenerContainer container, MessageConverter messageConverter) {
        Object messageListener = createMessageListener(container, messageConverter);
        Assert.state(messageListener != null, "Endpoint [" + this + "] must provide a non null message listener");
        if (this.retryTemplate != null) {
            if (messageListener instanceof AcknowledgingMessageListener) {
                messageListener = new RetryingAcknowledgingMessageListenerAdapter<>(
                        (AcknowledgingMessageListener<K, V>) messageListener, this.retryTemplate,
                        this.recoveryCallback);
            }
            else {
                messageListener = new RetryingMessageListenerAdapter<>((MessageListener<K, V>) messageListener,
                        this.retryTemplate, (RecoveryCallback<Object>) this.recoveryCallback);
            }
        }
        if (this.recordFilterStrategy != null) {
            if (messageListener instanceof AcknowledgingMessageListener) {
                messageListener = new FilteringAcknowledgingMessageListenerAdapter<>(
                        (AcknowledgingMessageListener<K, V>) messageListener, this.recordFilterStrategy,
                        this.ackDiscarded);
            }
            else if (messageListener instanceof MessageListener) {
                messageListener = new FilteringMessageListenerAdapter<>((MessageListener<K, V>) messageListener,
                        this.recordFilterStrategy);
            }
            else if (messageListener instanceof BatchAcknowledgingMessageListener) {
                messageListener = new FilteringBatchAcknowledgingMessageListenerAdapter<>(
                        (BatchAcknowledgingMessageListener<K, V>) messageListener, this.recordFilterStrategy,
                        this.ackDiscarded);
            }
            else if (messageListener instanceof BatchMessageListener) {
                messageListener = new FilteringBatchMessageListenerAdapter<>(
                        (BatchMessageListener<K, V>) messageListener, this.recordFilterStrategy);
            }
        }
        container.setupMessageListener(messageListener);
    }

如果retryTemplate不为null的话,会先判断是不是AcknowledgingMessageListener的子类,如果是则创建RetryingAcknowledgingMessageListenerAdapter,如果不是则创建RetryingMessageListenerAdapter

spring-kafka-1.2.3.RELEASE-sources.jar!/org/springframework/kafka/config/MethodKafkaListenerEndpoint.java

代码语言:javascript
复制
protected MessagingMessageListenerAdapter<K, V> createMessageListenerInstance(MessageConverter messageConverter) {
        if (isBatchListener()) {
            BatchMessagingMessageListenerAdapter<K, V> messageListener = new BatchMessagingMessageListenerAdapter<K, V>(
                    this.bean, this.method);
            if (messageConverter instanceof BatchMessageConverter) {
                messageListener.setBatchMessageConverter((BatchMessageConverter) messageConverter);
            }
            return messageListener;
        }
        else {
            RecordMessagingMessageListenerAdapter<K, V> messageListener =
                    new RecordMessagingMessageListenerAdapter<K, V>(this.bean, this.method);
            if (messageConverter instanceof RecordMessageConverter) {
                messageListener.setMessageConverter((RecordMessageConverter) messageConverter);
            }
            return messageListener;
        }
    }

其中RecordMessagingMessageListenerAdapter实现了AcknowledgingMessageListener接口

代码语言:javascript
复制
public class RecordMessagingMessageListenerAdapter<K, V> extends MessagingMessageListenerAdapter<K, V>
        implements MessageListener<K, V>, AcknowledgingMessageListener<K, V> {
        //......
}
本文参与 腾讯云自媒体分享计划,分享自微信公众号。
原始发表:2017-10-11,如有侵权请联系 cloudcommunity@tencent.com 删除

本文分享自 码匠的流水账 微信公众号,前往查看

如有侵权,请联系 cloudcommunity@tencent.com 删除。

本文参与 腾讯云自媒体分享计划  ,欢迎热爱写作的你一起参与!

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • AbstractRetryingMessageListenerAdapter
    • RetryingAcknowledgingMessageListenerAdapter
      • RetryingMessageListenerAdapter
      • 具体是哪种listener呢
      领券
      问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档