Spring Boot 2.7 Kafka消费者异常重试及max.poll.records疑问
我有一个Kafka监听器/消费者,消费拥有10个分区的topic“xyz”。在消息消费或处理过程中可能出现多种异常,我希望针对特定异常进行有限次数重试,已编写如下代码实现该需求:
@Bean public CommonErrorHandler customErrorHandler(){ DefaultErrorHandler defaultErrorHandler = new DefaultErrorHandler((rec,ex) -> log.error("Finished retries. Commiting offset"), new FixedBackOff(3000l, 5)); defaultErrorHandler.setCommitRecovered(true); defaultErrorHandler.addRetryableExceptions(SerializationException.class); return defaultErrorHandler; }
上述代码会在消费/处理失败时,以3秒间隔对同一消息重试5次,Spring Boot会自动将该Bean注入自动配置的监听器容器工厂。针对Spring Boot 2.7版本的问题解答如下:
问题1:是否可以配置泛型异常作为可重试异常,比如用Java的instanceOf判断异常类型来配置重试策略?
可以实现。Spring Kafka的DefaultErrorHandler支持通过自定义RetryPolicy实现基于类型匹配的重试判断,你可以在canRetry方法中用instanceof校验异常类型,再通过setRetryPolicy方法替换默认策略。示例代码如下:
RetryPolicy customRetryPolicy = new RetryPolicy() { @Override public boolean canRetry(RetryContext context) { Throwable throwable = context.getLastThrowable(); // 替换为你需要匹配的泛型异常类型 return throwable instanceof YourGenericException; } @Override public RetryContext open(RetryContext parent) { return new DefaultRetryContext(parent); } @Override public void close(RetryContext context) {} @Override public void registerThrowable(RetryContext context, Throwable throwable) {} }; defaultErrorHandler.setRetryPolicy(customRetryPolicy);
另外,也可以直接调用addRetryableExceptions添加父类异常,该父类的所有子类异常都会自动纳入重试范围,本质也是利用了instanceof的类型匹配逻辑。
问题2:是否可以针对特定异常进行无限重试?
可以实现。只需将FixedBackOff的重试次数设置为FixedBackOff.UNLIMITED_ATTEMPTS(对应值为-1),同时通过addRetryableExceptions指定需要无限重试的异常类型即可。示例修改如下:
DefaultErrorHandler defaultErrorHandler = new DefaultErrorHandler( (rec, ex) -> log.error("该兜底逻辑不会触发,因为是无限重试"), new FixedBackOff(3000l, FixedBackOff.UNLIMITED_ATTEMPTS) ); defaultErrorHandler.addRetryableExceptions(YourSpecificException.class);
需要注意:无限重试可能导致消费者长期卡在某条消息上,影响整体消费进度,建议搭配监控或告警机制使用。
问题3:关于Kafka消费者的“max.poll.records”属性:单消费者组的一个消费者消费10个分区且各分区均有消息的topic时,若该属性设为1,是每个分区拉取1条共10条消息,还是仅拉取1条消息?
当max.poll.records设为1时,单次poll操作会从每个分配的分区各拉取1条消息,也就是总共拉取10条消息。这个属性控制的是一次poll调用返回的最大记录总数,但Kafka消费者会在已分配的分区间均衡拉取数据,不会仅从单个分区获取。如果部分分区没有新消息,最终返回的记录数会小于10,但所有分区都有消息时会返回10条。
内容的提问来源于stack exchange,提问作者tinku

