すべてのプロダクト
Search
ドキュメントセンター

ApsaraMQ for RabbitMQ:クライアントの自動接続復旧の設定

最終更新日:Aug 22, 2026

ブローカーのアップグレード、再起動、ネットワークジッターにより、クライアントが切断されることがあります。自動回復は障害を検出し、手動での介入なしに再接続します。本ドキュメントでは、Java、Python、PHP の例を紹介します。

回復のトリガー

自動回復は、次の場合にトリガーされます。

  • I/O 例外がスローされた場合。

  • ソケットの読み取り操作がタイムアウトした場合。

  • サーバーハートビートが失われた場合。

Javaでの回復の設定

重要

amqp-client 4.0.0 以降では、接続復旧とトポロジー復旧はデフォルトで有効になっています。

以下のメソッドは、接続復旧とトポロジー復旧 (キュー、エクスチェンジ、バインディング、コンシューマー) を制御します。

amqp-client

説明

factory.setAutomaticRecoveryEnabled(boolean)

自動接続復旧を有効または無効にします。

factory.setNetworkRecoveryInterval(long)

再試行の間隔を設定します。デフォルトは 5 秒です。

factory.setTopologyRecoveryEnabled(boolean)

自動トポロジー復旧を有効または無効にします。

amqp-clientのサンプルコード

接続復旧とトポロジー復旧を有効にしたコンシューマークライアント:

ConnectionFactory factory = new ConnectionFactory();
// エンドポイント。ApsaraMQ for RabbitMQ コンソールの [インスタンスの詳細] ページでインスタンスのエンドポイントを取得できます。
factory.setHost("xxx.xxx.aliyuncs.com");
// ${instanceId} を ApsaraMQ for RabbitMQ インスタンスの ID に置き換えます。ApsaraMQ for RabbitMQ コンソールの [インスタンス] ページでインスタンス ID を取得できます。
factory.setCredentialsProvider(new AliyunCredentialsProvider("${instanceId}"));
// 仮想ホスト名。ApsaraMQ for RabbitMQ コンソールで仮想ホストが作成されていることを確認してください。
factory.setVirtualHost("${VhostName}");
// デフォルトポート。暗号化されていない接続にはポート 5672 を、暗号化された接続にはポート 5671 を使用します。
factory.setPort(5672);
// タイムアウト期間。ネットワーク環境に基づいて値を設定します。
factory.setConnectionTimeout(30 * 1000);
factory.setHandshakeTimeout(30 * 1000);
factory.setShutdownTimeout(0);
// 自動接続復旧を有効にするかどうかを指定します。
factory.setAutomaticRecoveryEnabled(true);
// 再試行間隔。値を 10 秒に設定します。
factory.setNetworkRecoveryInterval(10000);
// 自動トポロジー復旧を有効にするかどうかを指定します。
factory.setTopologyRecoveryEnabled(true);
Connection connection = factory.newConnection();

Pythonでの回復の設定

RabbitMQ が推奨する Python クライアントライブラリである Pika は、パラメーターによる自動回復をサポートしていません。代わりに、コールバック関数を使用して再接続ロジックを実装します。

Pika ベースの自動回復を備えたコンシューマークライアント:

# -*- coding: utf-8 -*-

import logging
import time
import pika

LOG_FORMAT = ('%(levelname) -10s %(asctime)s %(name) -30s %(funcName) '
              '-35s %(lineno) -5d: %(message)s')
LOGGER = logging.getLogger(__name__)
logging.basicConfig(level=logging.INFO, format=LOG_FORMAT)

class Consumer(object):

    def __init__(self, amqp_url, queue):
        self.should_reconnect = False

        self._connection = None
        self._channel = None
        self._closing = False
        self._url = amqp_url

        self._queue = queue

    def connect(self):
        '''
        接続を作成し、次のコールバックを設定します:
        on_open_callback:接続が作成されたときに呼び出されるコールバック。
        on_open_error_callback:接続の作成に失敗したときに呼び出されるコールバック。
        on_close_callback:接続が閉じられたときに呼び出されるコールバック。
        '''
        return pika.SelectConnection(
            parameters=pika.URLParameters(self._url),
            on_open_callback=self.on_connection_open,
            on_open_error_callback=self.on_connection_open_error,
            on_close_callback=self.on_connection_closed)

    def on_connection_open(self, _unused_connection):
        '''
        接続が作成されたときに呼び出されるコールバック。
        チャネルを作成し、次のコールバックを設定します:
        on_channel_open:チャネルが作成されたときに呼び出されるコールバック。
        '''
        self._connection.channel(on_open_callback=self.on_channel_open)

    def on_connection_open_error(self, _unused_connection, err):
        """
        接続の作成に失敗したときに呼び出されるコールバック。
        エラーメッセージを出力し、接続を再作成します。
        """
        LOGGER.error('Connection open failed: %s', err)
        self.reconnect()

    def on_connection_closed(self, _unused_connection, reason):
        """
        接続が閉じられたときに呼び出されるコールバック。
        次のシナリオが発生する可能性があります:
        1. 接続が正常に閉じられ、クライアントが終了します。
        2. クライアントが予期せず切断され、接続を再作成しようとします。
        """
        self._channel = None
        if self._closing:
            self._connection.ioloop.stop()
        else:
            LOGGER.warning('Connection closed, reconnect necessary: %s', reason)
            self.reconnect()

    def close_connection(self):
        """
        接続を閉じます。
        """
        if self._connection.is_closing or self._connection.is_closed:
            LOGGER.info('Connection is closing or already closed')
        else:
            LOGGER.info('Closing connection')
            self._connection.close()

    def reconnect(self):
        """
        self.should_reconnect パラメーターを True に設定し、I/O ループを停止します。
        """
        self.should_reconnect = True
        self.stop()

    def on_channel_open(self, channel):
        """
        チャネルが作成されたときに呼び出されるコールバック。
        コールバックを設定します。
        on_channel_closed:チャネルが閉じられたときに呼び出されるコールバック。
        キューからのメッセージ消費を開始します。
        """
        self._channel = channel
        self._channel.add_on_close_callback(self.on_channel_closed)
        self.start_consuming()

    def on_channel_closed(self, channel, reason):
        """
        チャネルが閉じられたときに呼び出されるコールバック。
        チャネル情報を出力し、接続を閉じます。
        """
        LOGGER.warning('Channel %i was closed: %s', channel, reason)
        self.close_connection()

    def start_consuming(self):
        """
        キューからのメッセージ消費を開始します。
        """
        LOGGER.info('start consuming...')
        self._channel.basic_consume(
            self._queue, self.on_message)

    def on_message(self, _unused_channel, basic_deliver, properties, body):
        """
        メッセージを消費し、確認応答 (ACK) を送信します。
        """
        LOGGER.info('Received message: %s', body.decode())
        # 消費ロジック。
        self._channel.basic_ack(basic_deliver.delivery_tag)

    def run(self):
        """
        接続を作成し、I/O ループを開始します。
        """
        self._connection = self.connect()
        self._connection.ioloop.start()

    def stop(self):
        """
        I/O ループを停止します。
        """
        if not self._closing:
            self._closing = True
            self._connection.ioloop.stop()
            LOGGER.info('Stopped')


class AutoRecoveryConsumer(object):

    def __init__(self, amqp_url, queue):
        self._amqp_url = amqp_url
        self._queue = queue
        self._consumer = Consumer(self._amqp_url, queue)

    def run(self):
        """
        KeyboardInterrupt 例外がスローされるまで while True ループを実行します。
        run() メソッドでは、I/O ループがリッスンするキューが開始され、メッセージが処理されます。このループにより、コンシューマーは継続的に実行し、ブローカーに自動的に再接続できます。
        """
        while True:
            try:
                self._consumer.run()
            except KeyboardInterrupt:
                self._consumer.stop()
                break
            self._maybe_reconnect()

    def _maybe_reconnect(self):
        """
        再接続が必要かどうかを判断します。2 つの連続した再接続の間隔は 1 秒です。
        """
        if self._consumer.should_reconnect:
            self._consumer.stop()
            time.sleep(1)
            self._consumer = Consumer(self._amqp_url, self._queue)


def main():
    username = 'MjoxODgwNzcwODY5MD****'
    password = 'NDAxREVDQzI2MjA0OT****'
    host = '1880770****.mq-amqp.cn-hangzhou-a.aliyuncs.com'
    port = 5672
    vhost = 'vhost_test'

    # amqp_url: amqp://:@:/
    amqp_url = 'amqp://%s:%s@%s:%i/%s' % (username, password, host, port, vhost)
    consumer = AutoRecoveryConsumer(amqp_url, 'QueueTest')
    consumer.run()


if __name__ == '__main__':
    main()

PHPでの回復の設定

php-amqplib は、AMQP 互換のメッセージキュー用の PHP ライブラリです。php-amqplib は、パラメーターによる自動回復をサポートしていません。再接続ロジックは手動で実装する必要があります。

php-amqplib ベースの自動回復を備えたコンシューマークライアント:

<?php

require_once __DIR__ . '/vendor/autoload.php';

use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Exception\AMQPRuntimeException;
use PhpAmqpLib\Exception\AMQPConnectionClosedException;

const ONE_SECOND = 1;


/**
 * 接続を作成します。
 */
function connect() {
    $host = '1880770****.mq-amqp.cn-hangzhou-a.aliyuncs.com';
    $username = 'MjoxODgwNzcwODY5MD****';
    $password = 'NDAxREVDQzI2MjA0OT****';
    $port = 5672;
    $vhost = 'vhost_test';
    return new AMQPStreamConnection($host, $port, $username, $password, $vhost);
}

/**
 * 接続を解放します。
 */
function cleanup_connection($connection) {
    try {
        if($connection !== null) {
            $connection->close();
        }
    } catch (\ErrorException $e) {
    }
}

$connection = null;

while(true){
    try {
        $connection = connect();
        start_consuming($connection);
    } catch (AMQPConnectionClosedException $e) {
        echo $e->getMessage() . PHP_EOL;
        cleanup_connection($connection);
        sleep(ONE_SECOND);
    } catch(AMQPRuntimeException $e) {
        echo $e->getMessage() . PHP_EOL;
        cleanup_connection($connection);
        sleep(ONE_SECOND);
    } catch(\RuntimeException $e) {
        echo 'Runtime exception ' . PHP_EOL;
        cleanup_connection($connection);
        sleep(ONE_SECOND);
    } catch(\ErrorException $e) {
        echo 'Error exception ' . PHP_EOL;
        cleanup_connection($connection);
        sleep(ONE_SECOND);
    }
}

/**
 * 消費を開始します。
 * @param AMQPStreamConnection $connection
 */
function start_consuming($connection) {
    $queue = 'queueTest';
    $consumerTag = 'consumer';
    $channel = $connection->channel();
    $channel->queue_declare($queue, false, true, false, false);
    $channel->basic_consume($queue, $consumerTag, false, false, false, false, 'process_message');
    while ($channel->is_consuming()) {
        $channel->wait();
    }
}


/**
 * メッセージを処理します。
 * @param \PhpAmqpLib\Message\AMQPMessage $message
 */
function process_message($message)
{
    // ビジネスロジックを処理します。
    echo "\n--------\n";
    echo $message->body;
    echo "\n--------\n";

    $message->ack();
}

制限事項

  • 接続障害の検出には時間がかかります。検出中のメッセージ損失を防ぐために、パブリッシャー確認を使用してください。

  • チャネル例外は自動回復をトリガーしません。これらのアプリケーションレベルのエラーは、コードで処理する必要があります。

  • 接続復旧ではチャネルは復元されません。

  • 接続が切断されると、その排他的キューは削除され、データはクリアされます。回復後、それらのキューからの消費は失敗します。

  • 接続上の各コンシューマーは、一意のコンシューマータグを持つ必要があります。コンシューマーがタグを共有している場合、1 つだけが回復されます。タグを省略すると、ブローカーが自動的にタグを割り当てます。