概要
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-1 | AMQP 1.0 |
|---|---|---|
| OASIS標準 | × | ○(2012年) |
| モデル | Exchange/Queue | Link/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)
実装・開発のポイント
主要ブローカーと対応ライブラリ
| ブローカー | プロトコル | 特徴 |
|---|---|---|
| RabbitMQ | AMQP 0-9-1(+1.0プラグイン) | 最も普及、管理画面充実 |
| Azure Service Bus | AMQP 1.0 | Azureとの統合 |
| ActiveMQ | AMQP 0-9-1/1.0, MQTT等 | 多プロトコル対応 |
| Qpid | AMQP 1.0 | Apache財団、Redhat系 |
| クライアントライブラリ | 言語 | バージョン |
|---|---|---|
| pika | Python | 0-9-1 |
| aio-pika | Python(非同期) | 0-9-1 |
| amqplib | C | 0-9-1 |
| uamqp | Python | 1.0 |
| rhea | JavaScript | 1.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
)
)
他技術との比較
| 項目 | AMQP | MQTT | Kafka | ZeroMQ |
|---|---|---|---|---|
| メッセージ保証 | at-least-once | 0/1/2 | at-least-once | best-effort |
| ブローカー必要 | 必要 | 必要 | 必要 | 不要 |
| ルーティング柔軟性 | ◎ | △ | △ | ◎ |
| スループット | 中 | 中 | ◎ | ◎ |
| 標準化 | OASIS標準 | OASIS標準 | — | — |
| 組み込み適合 | △ | ◎ | × | ○ |
MQTTは超軽量デバイスとのエッジ通信に向き、AMQPはバックエンドの信頼性重視メッセージングに向きます。Pub/Subモデルの実装として、AMQPはTopicエクスチェンジ、MQTTはブローカーのトピック機能を使います。エンタープライズIoTでは、エッジ側MQTT・バックエンド側AMQPのハイブリッド構成が現実的な選択です。