【问题标题】:Reliable Messaging with RabbitMQ使用 RabbitMQ 进行可靠的消息传递
【发布时间】:2014-11-27 01:01:23
【问题描述】:

我有一个通过 RabbitMQ 发送 AMQP 消息的应用程序。消息发送是在 http 请求上触发的。最近我注意到有些消息似乎丢失了(就像从未发送过一样)。我还注意到服务器管理的频道列表正在稳步增加。我纠正的第一件事是在不再需要通道后关闭它们。但是,我仍然不确定我的代码结构是否正确以确保交付。下面是两段代码;第一个是管理连接的单例部分(不会在每次调用时重新创建),第二个是发送代码。任何建议/指导将不胜感激。

@Service
public class PersistentConnection {
    private static Connection myConnection = null;
    private Boolean blocked = false;

    @Autowired ApplicationConfiguration applicationConfiguration;
    @Autowired ConfigurationService configurationService;

    @PostConstruct
    private void init() {
    }

    @PreDestroy
    private void destroy() {
        try {
            myConnection.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    public Connection getConnection( ) {
        if (myConnection == null) {
            start();
        }
        else if (!myConnection.isOpen()) {
            log.warn("AMQP Connection closed.  Attempting to start.");
            start();
        }
        return myConnection;
    }


    private void start() {
        log.debug("Building AMQP Connection");

        ConnectionFactory factory = new ConnectionFactory();
        String ipAddress = applicationConfiguration.getAMQPHost();
        String password = applicationConfiguration.getAMQPUser();
        String user = applicationConfiguration.getAMQPPassword();
        String virtualHost = applicationConfiguration.getAMQPVirtualHost();
        String port = applicationConfiguration.getAMQPPort();

        try {
            factory.setUsername(user);
            factory.setPassword(password);
            factory.setVirtualHost(virtualHost);
            factory.setPort(Integer.parseInt(port));
            factory.setHost(ipAddress);
            myConnection = factory.newConnection();
        }
        catch (Exception e) {
            e.printStackTrace();
        }

        myConnection.addBlockedListener(new BlockedListener() {
            public void handleBlocked(String reason) throws IOException {
                // Connection is now blocked
                blocked = true;
            }

            public void handleUnblocked() throws IOException {
                // Connection is now unblocked
                blocked = false;
            }
        });
    }

    public Boolean isBlocked() {
        return blocked;
    }
}

/*
 * Sends ADT message to AMQP server.
 */
private void send(String routingKey, String message) throws Exception { 
    String exchange = applicationConfiguration.getAMQPExchange();  
    String exchangeType = applicationConfiguration.getAMQPExchangeType();

    Connection connection = myConnection.getConnection();
    Channel channel = connection.createChannel();
    channel.exchangeDeclare(exchange, exchangeType);
    channel.basicPublish(exchange, routingKey, null, message.getBytes());

    // Close the channel if it is no longer needed in this thread
    channel.close();
}

【问题讨论】:

  • 可以从更多线程调用getConnection( ) 吗?如果是,则代码不是线程安全的。
  • 一个 http 请求进来,最终调用 getConnection() 调用。所以我想它可以。我想通过使它成为一个单例来解决这个问题。改进代码的最佳方法是什么?

标签: rabbitmq amqp


【解决方案1】:

试试这个代码:

@Service
public class PersistentConnection {
    private Connection myConnection = null;
    private Boolean blocked = false;

    @Autowired ApplicationConfiguration applicationConfiguration;
    @Autowired ConfigurationService configurationService;

    @PostConstruct
    private void init() {
      start(); /// In this way you can initthe connection and you are sure it is called only one time.
    }

    @PreDestroy
    private void destroy() {
        try {
            myConnection.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    public Connection getConnection( ) {
        return myConnection;
    }


    private void start() {
        log.debug("Building AMQP Connection");

        ConnectionFactory factory = new ConnectionFactory();
        String ipAddress = applicationConfiguration.getAMQPHost();
        String password = applicationConfiguration.getAMQPUser();
        String user = applicationConfiguration.getAMQPPassword();
        String virtualHost = applicationConfiguration.getAMQPVirtualHost();
        String port = applicationConfiguration.getAMQPPort();

        try {
            factory.setUsername(user);
            factory.setPassword(password);
            factory.setVirtualHost(virtualHost);
            factory.setPort(Integer.parseInt(port));
            factory.setHost(ipAddress);
            myConnection = factory.newConnection();
        }
        catch (Exception e) {
            e.printStackTrace();
        }

        myConnection.addBlockedListener(new BlockedListener() {
            public void handleBlocked(String reason) throws IOException {
                // Connection is now blocked
                blocked = true;
            }

            public void handleUnblocked() throws IOException {
                // Connection is now unblocked
                blocked = false;
            }
        });
    }

    public Boolean isBlocked() {
        return blocked;
    }
}

/*
 * Sends ADT message to AMQP server.
 */
private void send(String routingKey, String message) throws Exception { 
    String exchange = applicationConfiguration.getAMQPExchange();  
    String exchangeType = applicationConfiguration.getAMQPExchangeType();

    Connection connection = myConnection.getConnection();
    if (connection!=null){
    Channel channel = connection.createChannel();
    try{
    channel.exchangeDeclare(exchange, exchangeType);
    channel.basicPublish(exchange, routingKey, null, message.getBytes());
    } finally{
      // Close the channel if it is no longer needed in this thread
       channel.close();
   }

} }

这样就够了,系统启动时你已经和rabbitmq建立了连接。

如果你是一个懒惰的单身人士,代码只是有点不同。

我建议不要使用isOpen()方法,请阅读here

打开

boolean isOpen() 判断组件当前是否打开。 如果我们当前正在关闭,将返回 false。检查此方法 由于竞争条件,应该仅供参考 - 状态 通话后可以更改。相反,只需执行并尝试捕捉 ShutdownSignalException 和 IOException 返回:组件时为 true 是打开的,否则为假

编辑**

问题一:

你要找的是 HA 客户端。

RabbitMQ java客户端默认不支持此功能,因为3.3.0版本只支持重连,阅读this

...允许基于 Java 的客户端在网络后自动重新连接 失败。如果您想确定您的消息,您必须创建一个 强大的客户端能够抵抗所有失败。

通常你应该考虑失败,例如: 消息发布过程中出现错误怎么办?

在您的情况下,您只是丢失了消息,您应该手动重新排队消息。

问题 2:

我不知道你的代码,但 connection == null 不应该发生,因为这个过程首先被调用:

@PostConstruct
    private void init() {
      start(); /// In this way you can initthe connection and you are sure it is called only one time.
    }

无论如何你都可以提出异常,问题是: 我与我试图发送的消息有什么关系?

见问题 1

我想建议阅读有关 HA 的更多信息,例如:

用rabbitmq创建一个可靠的系统并不复杂,但你应该知道一些基本概念。

无论如何.. 让我知道!

【讨论】:

  • 非常感谢您。我今晚会处理它,明天加载代码并更新你。
  • 我刚开始研究这个,我有几个问题。首先,如果 RabbitMQ 服务器重新启动(例如)会发生什么 - 我假设连接将失败并且在服务器启动时不会自动重新连接?有没有办法自动重新连接?另外,我在想如果 connection==null ,发送方法可能应该抛出一个异常 - 这有意义吗?
  • 这真的很有帮助。我已经加载了新代码,它似乎运行良好。我还通读了参考资料——找到了你的书 :)。我不认为我的应用程序真的需要 Lyra,但我会记住它以备将来使用。再次感谢您。
  • @skyman 如果可靠性是一个问题,请检查 Lyra。尝试推出自己的 HA 客户端看似很难,而且您要么没有涵盖所有可能的场景,要么最终实现了 Lyra 之类的东西:)
  • @Jonathan 感谢您的评论。我似乎仍然有一些小问题 - 会看看 Lyra
猜你喜欢
  • 1970-01-01
  • 2019-12-27
  • 1970-01-01
  • 1970-01-01
  • 2012-03-23
  • 1970-01-01
  • 2013-02-17
  • 2016-06-14
  • 2019-06-03
相关资源
最近更新 更多