Spring Boot中定义多Pulsar Producer,简化消息发送代码方案咨询
简化Spring Boot Pulsar多主题消息发送的方案
核心思路
通过自定义注解绑定消息类型与对应Pulsar主题、JSON Schema,再封装一个通用发送器自动完成Schema配置和主题匹配,最终业务代码仅需传入消息实例即可完成发送,彻底消除重复的setSchema和硬编码主题名操作。
具体实现步骤
1. 定义消息绑定注解
创建运行时注解,用来标记每个消息类对应的目标主题和Schema类型:
@Target(ElementType.TYPE) @Retention(RetentionPolicy.RUNTIME) public @interface PulsarMessageBinding { // 绑定的Pulsar主题名 String topic(); // 默认使用JSON Schema,可扩展支持其他Schema类型 Class<? extends Schema<?>> schema() default JSONSchema.class; }
2. 为消息实体绑定注解
给每个业务消息类添加上述注解,明确其对应的主题:
@PulsarMessageBinding(topic = "user-register-topic") public class UserRegisterEvent { private String userId; private String email; // getter/setter 省略 } @PulsarMessageBinding(topic = "order-create-topic") public class OrderCreateEvent { private String orderId; private BigDecimal amount; // getter/setter 省略 }
3. 封装通用Pulsar发送器
基于Spring的PulsarTemplate,结合注解解析实现自动适配:
@Component public class GenericPulsarSender { private final PulsarTemplate<Object> pulsarTemplate; private final Map<Class<?>, PulsarMessageBinding> bindingCache = new ConcurrentHashMap<>(); private final Map<Class<?>, Schema<Object>> schemaCache = new ConcurrentHashMap<>(); public GenericPulsarSender(PulsarTemplate<Object> pulsarTemplate, ApplicationContext context) { this.pulsarTemplate = pulsarTemplate; // 初始化时扫描所有带注解的消息类,缓存绑定关系 Map<String, Object> annotatedBeans = context.getBeansWithAnnotation(PulsarMessageBinding.class); annotatedBeans.values().forEach(bean -> { Class<?> clazz = bean.getClass(); PulsarMessageBinding binding = clazz.getAnnotation(PulsarMessageBinding.class); bindingCache.put(clazz, binding); // 预创建Schema实例缓存,提升性能 schemaCache.put(clazz, (Schema<Object>) JSONSchema.of(clazz)); }); } public void send(Object message) { Class<?> msgClass = message.getClass(); PulsarMessageBinding binding = bindingCache.get(msgClass); if (binding == null) { throw new IllegalArgumentException("未找到消息类[" + msgClass.getName() + "]对应的Pulsar绑定配置"); } // 自动应用缓存的Schema和主题 pulsarTemplate.setSchema(schemaCache.get(msgClass)); pulsarTemplate.send(binding.topic(), message); } // 扩展异步发送方法 public CompletableFuture<MessageId> sendAsync(Object message) { Class<?> msgClass = message.getClass(); PulsarMessageBinding binding = bindingCache.get(msgClass); if (binding == null) { throw new IllegalArgumentException("未找到消息类[" + msgClass.getName() + "]对应的Pulsar绑定配置"); } pulsarTemplate.setSchema(schemaCache.get(msgClass)); return pulsarTemplate.sendAsync(binding.topic(), message); } }
4. 业务层简化调用
业务代码只需注入通用发送器,直接传入消息实例即可完成发送:
@Service public class UserService { private final GenericPulsarSender pulsarSender; public UserService(GenericPulsarSender pulsarSender) { this.pulsarSender = pulsarSender; } public void processUserRegister(UserRegisterEvent event) { // 业务逻辑处理... pulsarSender.send(event); // 仅需传入消息实例,无需手动配置Schema和主题 } }
额外优化建议
- 增加全局异常处理,捕获发送失败异常并做降级或重试处理
- 支持自定义Schema的初始化逻辑,比如针对特殊字段的JSON序列化配置
- 利用Spring的
BeanPostProcessor替代初始化时的扫描,实现更灵活的注解处理
内容的提问来源于stack exchange,提问作者Gootik
相关产品推荐
相关产品推荐

