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

如何结合Alpakka使用Source.queue?复用JMS连接的生产者实现

实现可复用的JMS队列生产者Actor

完全理解你的需求——不想每次发消息都重新创建JMS连接,而是用一个Actor维持住连接,所有消息都走同一个复用的流程。下面是完整的实现方案,基于你给出的初始化代码扩展:

核心思路

我们用Akka Actor封装JMS生产流:

  • Actor启动时初始化JMS Sink和一个Source.queue(作为消息入口)
  • 将Source.queue与JMS Sink连接成持续运行的流,这样JMS连接会一直保持打开状态
  • 外部通过给Actor发消息的方式提交要发送的内容,Actor把消息推入Source.queue,复用已有的流和连接

完整Actor实现

import akka.actor.AbstractActor;
import akka.actor.ActorRef;
import akka.actor.Props;
import akka.stream.javadsl.JmsProducer;
import akka.stream.javadsl.Sink;
import akka.stream.javadsl.Source;
import akka.stream.OverflowStrategy;
import java.util.concurrent.CompletionStage;

// 定义Actor接收的消息类型
static class SendMessage {
    public final String message;
    public SendMessage(String message) {
        this.message = message;
    }
}

static class JmsProducerActor extends AbstractActor {
    private final Source.QueueWithComplete<String> messageQueue;
    private final CompletionStage<Done> streamCompletion;

    public static Props props(JmsProducerSettings jmsSettings) {
        return Props.create(JmsProducerActor.class, jmsSettings);
    }

    public JmsProducerActor(JmsProducerSettings jmsSettings) {
        // 初始化JMS Sink(和你给出的代码一致)
        Sink<String, CompletionStage<Done>> jmsSink = JmsProducer.textSink(jmsSettings);
        
        // 创建消息入口队列,配置缓冲区和背压策略
        this.messageQueue = Source.<String>queue(Integer.MAX_VALUE, OverflowStrategy.backpressure())
                .toMat(jmsSink, Source::preMaterialize)
                .run(getContext().getSystem());
        
        // 保存流的完成阶段,用于Actor停止时清理
        this.streamCompletion = messageQueue.completionStage();
    }

    @Override
    public Receive createReceive() {
        return receiveBuilder()
                // 处理发送消息的请求
                .match(SendMessage.class, msg -> {
                    messageQueue.offer(msg.message)
                            .whenComplete((queueResult, ex) -> {
                                if (ex != null) {
                                    // 处理入队失败的情况,比如流已经终止
                                    getSelf().tell(new MessageSendFailed(ex), getSender());
                                } else {
                                    // 根据队列结果处理,比如返回成功确认
                                    switch (queueResult) {
                                        case Dropped:
                                            getSender().tell(new MessageDropped(), getSelf());
                                            break;
                                        case QueueClosed:
                                            getSender().tell(new ProducerStopped(), getSelf());
                                            break;
                                        case Enqueued:
                                            getSender().tell(new MessageSent(), getSelf());
                                            break;
                                    }
                                }
                            });
                })
                // 处理Actor停止的逻辑,关闭队列并等待流完成
                .matchEquals("stop", msg -> {
                    messageQueue.complete();
                    streamCompletion.thenRun(() -> getContext().stop(getSelf()));
                })
                .build();
    }

    // 可选的响应消息类型,用于通知发送端结果
    static class MessageSent {}
    static class MessageDropped {}
    static class ProducerStopped {}
    static class MessageSendFailed {
        public final Throwable cause;
        public MessageSendFailed(Throwable cause) { this.cause = cause; }
    }
}

关键细节说明

  • 连接复用:JmsProducer内部会管理JMS连接、会话和生产者,只要流处于运行状态,这些资源就会保持打开,不会为每条消息重新创建。只有当Actor停止、流关闭时,资源才会被释放。
  • Source.queue配置:你用了Integer.MAX_VALUE作为缓冲区大小和OverflowStrategy.backpressure,这个配置适合消息发送速度不超过JMS处理速度的场景——如果缓冲区满了,offer操作会等待缓冲区有空位,避免消息丢失。如果需要非阻塞的行为,可以换成OverflowStrategy.dropNew之类的策略。
  • 错误处理:示例中添加了对入队结果的处理,你可以根据业务需求调整——比如当流因JMS连接失败终止时,Actor可以选择重启、通知监控系统等。

使用示例

// 初始化JMS设置(你的connectionProvider和queueName在这里传入)
JmsProducerSettings jmsSettings = JmsProducerSettings.create(connectionProvider.get())
        .withQueue(queueName);

// 创建JMS生产者Actor
ActorRef jmsProducer = system.actorOf(JmsProducerActor.props(jmsSettings), "jms-producer");

// 发送消息
jmsProducer.tell(new SendMessage("Hello JMS!"), ActorRef.noSender());

// 停止Actor(清理资源)
jmsProducer.tell("stop", ActorRef.noSender());

内容的提问来源于stack exchange,提问作者Kingpin2k

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:21:52