如何复用Kafka消费者中多处重复的异常处理块,避免代码冗余
适配Kafka @StreamListener的原生全局异常处理方案
Spring Cloud Stream 本身提供了和 @ControllerAdvice 逻辑完全对齐的异常处理能力,不需要在业务方法里手写重复的 try-catch 块:
你可以直接定义全局异常处理器,监听 Spring Cloud Stream 内置的全局错误通道 errorChannel,所有 @StreamListener 方法抛出的未捕获异常都会流转到这个通道统一处理,示例代码如下:
@Component public class GlobalStreamErrorHandler { @ServiceActivator(inputChannel = "errorChannel") public void handleStreamError(ErrorMessage errorMessage) { Throwable ex = errorMessage.getPayload(); if (ex instanceof CustomException) { log.error("Custom exception occurred"); } else if (ex instanceof MongoException) { log.error("Mongo exception occurred"); throw (MongoException) ex; } else if (ex instanceof ResourceNotFoundException) { log.error("ResourceNotFound exception occurred"); } else { log.error("Something totally weird happened"); throw new RuntimeException(ex); } } }
这个方案是 Kafka 消费场景的官方标准实现,和 @ControllerAdvice 的使用体验完全一致,无需侵入业务代码,是最优选择。
通用代码复用方案(不依赖Spring Cloud Stream特性)
如果需要脱离 Spring Cloud Stream 特性使用,或者需要更灵活的自定义逻辑,可以选择以下方案:
方案1:函数式接口封装模板方法
把异常处理逻辑封装为通用工具类,接收函数式接口作为业务逻辑入参,改造成本极低:
public class GlobalExceptionUtil { public static void runWithExceptionHandle(Runnable businessLogic) { try { businessLogic.run(); log.info("Finished processing business logic"); } catch (CustomException ex) { log.error("Custom exception occurred"); } catch (MongoException ex) { log.error("Mongo exception occurred"); throw ex; } catch (ResourceNotFoundException ex) { log.error("ResourceNotFound exception occurred"); } catch (Exception ex) { log.error("Something totally weird happened"); throw ex; } } }
使用方式如下,只需要把原来的业务逻辑放入Lambda表达式即可:
public class A { public void funcA(Input input){ GlobalExceptionUtil.runWithExceptionHandle(() -> { randomService.randomFunction(input); }); } }
方案2:AOP切面无侵入处理
如果不需要修改任何业务代码,可以用Spring AOP环绕增强所有 @StreamListener 标注的方法,统一处理异常:
@Aspect @Component public class StreamListenerExceptionAspect { @Around("@annotation(org.springframework.cloud.stream.annotation.StreamListener)") public Object aroundStreamListener(ProceedingJoinPoint pjp) throws Throwable { try { return pjp.proceed(); } catch (CustomException ex) { log.error("Custom exception occurred"); return null; } catch (MongoException ex) { log.error("Mongo exception occurred"); throw ex; } catch (ResourceNotFoundException ex) { log.error("ResourceNotFound exception occurred"); return null; } catch (Exception ex) { log.error("Something totally weird happened"); throw ex; } } }
方案3:父类模板方法继承
如果所有消费者结构高度统一,可以把异常处理逻辑放到抽象父类,子类只需要实现业务逻辑即可:
public abstract class AbstractKafkaConsumer<T> { public void process(T input) { try { doBusinessProcess(input); log.info("Finished processing business logic"); } catch (CustomException ex) { log.error("Custom exception occurred"); } catch (MongoException ex) { log.error("Mongo exception occurred"); throw ex; } catch (ResourceNotFoundException ex) { log.error("ResourceNotFound exception occurred"); } catch (Exception ex) { log.error("Something totally weird happened"); throw ex; } } protected abstract void doBusinessProcess(T input); }
子类实现示例:
public class A extends AbstractKafkaConsumer<Input> { @Autowired private RandomService randomService; @Override protected void doBusinessProcess(Input input) { randomService.randomFunction(input); } @StreamListener("input-channel") public void funcA(Input input) { process(input); } }
选型建议
- 优先选择Spring Cloud Stream原生全局错误通道方案,是官方针对该场景的标准实现
- 需要脱离Spring Cloud Stream使用的场景优先选择函数式接口封装方案,简单易维护
- 需要完全无侵入业务代码的场景选择AOP方案
- 所有消费者逻辑高度同质化的场景可以选择父类模板方法方案
内容的提问来源于stack exchange,提问作者Kaushik
相关产品推荐
相关产品推荐

