IoTプロトコル

AMQP

信頼性重視のメッセージングプロトコル。

概要

AMQP(Advanced Message Queuing Protocol)は、メッセージ指向ミドルウェア(MOM:Message-Oriented Middleware)のための開放的な標準プロトコルです。2006年にJPMorganが中心となり、Cisco、RedHat、Microsoft、VMwareなど多くの企業が参加して策定されました。2012年にはOASIS標準(バージョン1.0)として採択されています。

MQTTが軽量・省電力を優先するのに対し、AMQPはメッセージの信頼性・柔軟なルーティング・相互運用性を重視して設計されています。金融取引システム、エンタープライズアプリケーション統合(EAI)、大規模IoTバックエンドの中継層など、「絶対にメッセージを落としてはいけない」用途に採用されます。

主なAMQPブローカーにはRabbitMQ、Azure Service Bus、ActiveMQなどがあります。IoT文脈では、デバイスから収集したデータをクラウドのメッセージキューで処理するバックエンド統合層として使われることが多く、エッジデバイス側はMQTTやHTTPで通信し、ゲートウェイでAMQPに変換するアーキテクチャが一般的です。


歴史・背景

2000年代初頭、企業間のシステム統合にはIBM MQ(旧MQSeries)やTIBCO Rendezvousといった独自プロトコルのメッセージングミドルウェアが使われていました。これらは高機能ですが、ベンダーロックインが激しく、異なるメーカーのブローカー間で相互運用できないという根本的な問題がありました。

JPMorganのJohn O’Hara氏はこの状況を改善するため、2003年に標準化プロジェクトを立ち上げました。2006年にAMQP 0-9-1が公開され、2012年にOASISがAMQP 1.0を標準化しました。なお、RabbitMQなど多くの実装が依然として0-9-1を使っており、AMQP 1.0への移行は緩やかに進んでいます。

IoT分野では、Azure IoT HubがAMQP 1.0をネイティブサポートしており、高スループット・低レイテンシが求められる産業IoT(IIoT)バックエンドでの採用が進んでいます。


技術仕様

AMQPモデルの構成要素

AMQP 0-9-1(RabbitMQ等で広く使われるバージョン)のモデルを説明します。

Producer  →  [Exchange]  →  [Binding]  →  [Queue]  →  Consumer
                 ↑                            ↑
            ルーティングキー              メッセージ格納
構成要素説明
Producer(発行者)メッセージを送信するアプリケーション
Exchange(交換機)メッセージを受け取りキューに振り分ける
Queue(キュー)メッセージを格納するバッファ
Binding(バインディング)ExchangeとQueueの接続ルール
Consumer(消費者)キューからメッセージを受け取るアプリケーション

Exchangeの種類

タイプ振り分けルール用途
Directルーティングキー完全一致特定キューへの配信
Topicワイルドカード一致(*, #パターンマッチング
Fanoutすべてのバインドキューに配信ブロードキャスト
Headersメッセージヘッダー属性で振り分け複雑なルーティング
# Topicエクスチェンジの例
sensors.factory1.temperature → factory1_tempキューにルーティング
sensors.*.temperature        → 全工場の温度データ
sensors.factory1.#           → 工場1のすべてのセンサー

メッセージ確認(Acknowledgement)

AMQPの最大の特徴がメッセージ配信確認です。

Consumer                          Broker
  |                                  |
  |-- Basic.Consume ---------------->|  購読開始
  |<-- Basic.Deliver (tag=1) --------|  メッセージ配信
  |  [メッセージ処理]                |
  |-- Basic.Ack (tag=1) ------------>|  処理完了を通知
  |  ← ブローカーがキューから削除    |

ACKを送らない場合、ブローカーはメッセージをキューに保持し続け、接続が切れると別のコンシューマーに再配信します。これにより「少なくとも1回(at-least-once)」の配信を保証します。

QoSとプリフェッチ

# コンシューマーが一度に受け取るメッセージ数を制限
channel.basic_qos(prefetch_count=10)

プリフェッチ制御により、処理が遅いコンシューマーがメッセージを大量に受け取って詰まることを防ぎます。

AMQP 1.0 vs 0-9-1

項目AMQP 0-9-1AMQP 1.0
OASIS標準×○(2012年)
モデルExchange/QueueLink/Node
RabbitMQ対応○(デフォルト)プラグインで対応
Azure IoT対応○(ネイティブ)
複雑度

動作原理

接続確立

AMQPはTCPポート5672(TLSは5671)を使います。

Client                              Broker
  |-- TCP Connect ----------------->|
  |-- AMQP Header (protocol ver) -->|
  |<-- AMQP Header -----------------|
  |-- Connection.Start-OK ---------->|  SASL認証
  |<-- Connection.Tune -------------|  フレームサイズネゴシエート
  |-- Connection.Tune-OK ----------->|
  |-- Connection.Open (vhost) ------>|
  |<-- Connection.Open-OK ----------|
  |-- Channel.Open ----------------->|
  |<-- Channel.Open-OK --------------|
  |  [Exchange宣言・Queue宣言・Bind等]
  |-- Basic.Publish ----------------->|  メッセージ送信

メッセージの永続化

import pika

# キューを永続化(ブローカー再起動でも失わない)
channel.queue_declare(queue='sensor_data', durable=True)

# メッセージを永続化
properties = pika.BasicProperties(
    delivery_mode=pika.spec.PERSISTENT_DELIVERY_MODE
)
channel.basic_publish(
    exchange='',
    routing_key='sensor_data',
    body=json.dumps(data),
    properties=properties
)

Dead Letter Queue(DLQ)

処理に失敗したメッセージを別のキューに移動させる仕組みです。

# メインキューにDLX(Dead Letter Exchange)を設定
args = {
    'x-dead-letter-exchange': 'dlx_exchange',
    'x-dead-letter-routing-key': 'failed_messages',
    'x-message-ttl': 30000  # 30秒でDLQへ
}
channel.queue_declare(queue='sensor_data', durable=True, arguments=args)

用途・ユースケース

IoTデータパイプライン

多数のIoTデバイスからのデータを信頼性高くバックエンドに収集します。

[ESP32 × 1000台]
    │ (MQTT)

[MQTTブローカー]
    │ (AMQP Bridge)

[RabbitMQ/Azure Service Bus]

    ├→ [時系列DB書き込みワーカー]
    ├→ [アラート処理ワーカー]
    └→ [機械学習パイプライン]

工場ライン管理(IIoT)

生産ラインの各機械がAMQPで品質データを送信し、品質管理システムが非同期に処理します。メッセージが消えない永続キューにより、システム停止中のデータも保全されます。

# 生産ラインデータの送信(Python)
import pika
import json
from datetime import datetime

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='mes-server.factory.local',
        port=5672,
        credentials=pika.PlainCredentials('iot_device', 'secret')
    )
)
channel = connection.channel()

# 品質データを送信
data = {
    "machine_id": "CNC_LATHE_003",
    "timestamp": datetime.utcnow().isoformat(),
    "part_id": "P-2025-001234",
    "diameter_mm": 25.012,
    "tolerance_ok": True
}

channel.basic_publish(
    exchange='production.topic',
    routing_key='quality.cnc.lathe',
    body=json.dumps(data),
    properties=pika.BasicProperties(delivery_mode=2)
)
connection.close()

Azure IoT Hubとの統合

# Azure IoT Hub への AMQP 1.0 接続(Pythonの uamqp ライブラリ)
import uamqp
from uamqp import authentication

connection_string = "HostName=myhub.azure-devices.net;..."
uri = "amqps://myhub.azure-devices.net/devices/device_01/messages/events"

sas_auth = authentication.SASTokenAuth.from_connection_string(
    connection_string
)

with uamqp.SendClient(uri, auth=sas_auth) as send_client:
    msg = uamqp.Message(b'{"temperature": 25.3}')
    send_client.send_message(msg)

実装・開発のポイント

主要ブローカーと対応ライブラリ

ブローカープロトコル特徴
RabbitMQAMQP 0-9-1(+1.0プラグイン)最も普及、管理画面充実
Azure Service BusAMQP 1.0Azureとの統合
ActiveMQAMQP 0-9-1/1.0, MQTT等多プロトコル対応
QpidAMQP 1.0Apache財団、Redhat系
クライアントライブラリ言語バージョン
pikaPython0-9-1
aio-pikaPython(非同期)0-9-1
amqplibC0-9-1
uamqpPython1.0
rheaJavaScript1.0

コンシューマーの実装パターン

import pika
import json

def process_sensor_data(ch, method, properties, body):
    """センサーデータの処理"""
    try:
        data = json.loads(body)
        
        # データをDBに保存
        save_to_database(data)
        
        # 成功したらACK
        ch.basic_ack(delivery_tag=method.delivery_tag)
        
    except Exception as e:
        # 失敗したらNACK(再キュー)
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True  # Falseにすると DLQ へ
        )

channel.basic_qos(prefetch_count=5)
channel.basic_consume(
    queue='sensor_data',
    on_message_callback=process_sensor_data
)
channel.start_consuming()

接続断時のハートビート

# ハートビート設定(NAT/FWのタイムアウトを防ぐ)
connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host='broker.example.com',
        heartbeat=30,        # 30秒ごとにハートビート
        blocked_connection_timeout=300
    )
)

他技術との比較

項目AMQPMQTTKafkaZeroMQ
メッセージ保証at-least-once0/1/2at-least-oncebest-effort
ブローカー必要必要必要必要不要
ルーティング柔軟性
スループット
標準化OASIS標準OASIS標準
組み込み適合×

MQTTは超軽量デバイスとのエッジ通信に向き、AMQPはバックエンドの信頼性重視メッセージングに向きます。Pub/Subモデルの実装として、AMQPはTopicエクスチェンジ、MQTTはブローカーのトピック機能を使います。エンタープライズIoTでは、エッジ側MQTT・バックエンド側AMQPのハイブリッド構成が現実的な選択です。

関連用語

参考リンク