You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 00:00:08