概要
Pub/Sub(Publish/Subscribe、パブリッシュ/サブスクライブ)は、情報の発行者(Publisher)と購読者(Subscriber)を直接接続せずに分離するメッセージング通信パターンです。両者はブローカー(仲介者)またはメッセージバスを通じて間接的にやり取りし、互いの存在を知る必要がありません。
この「疎結合(Loose Coupling)」が最大の特徴であり、以下の3つの次元で分離を実現します。
- 空間的分離:発行者と購読者がお互いのIPアドレスを知らなくてよい
- 時間的分離:発行者がメッセージを送った時点で購読者がオンラインである必要がない
- 同期的分離:発行者はメッセージを送った後、応答を待つ必要がない
IoT分野では、数千〜数万台のセンサーデバイスからのデータを複数のバックエンドサービスが処理する構成においてPub/Subが基本アーキテクチャとなっています。MQTTはPub/Subの最も有名な実装の一つであり、他にもAMQP・OPC UA Pub/Sub・Apache Kafkaなどが採用しています。
歴史・背景
Pub/Subパターンの起源は1987年のTIBCO Rendezvous(金融市場データ配信)にさかのぼります。その後、エンタープライズメッセージングの世界でJMS(Java Message Service)が普及し、Pub/Subは標準的なメッセージングパターンの一つとして認知されました。
1990年代後半にはインターネット規模での情報配信の需要が高まり、さまざまなPub/Sub実装が生まれました。Usenetのニュースグループ、RSSフィード(2000年代)もPub/Subの一形態です。
IoT分野への本格展開は2008年頃から始まり、MQTTのOASIS標準化(2014年)、AWS IoT/Google Cloud IoTの登場(2015〜2016年)によって爆発的に普及しました。現在ではGoogle Cloud Pub/Sub、AWS SNS/SQS、Azure Event HubsなどすべてのメジャークラウドがPub/Subサービスを提供しています。
技術仕様
基本モデル
Publisher A ─┐ ┌─ Subscriber X
Publisher B ─┼─→ [ブローカー] ─→ ├─ Subscriber Y
Publisher C ─┘ (メッセージ └─ Subscriber Z
フィルタリング)
トピックベースとコンテンツベース
Pub/Subのフィルタリング方式は大きく2種類に分かれます。
トピックベース(Topic-based):
メッセージをトピック(チャンネル)に分類し、購読者はトピックを指定して購読します。MQTTが代表例です。
トピック: sensors/factory1/line1/temperature
トピック: sensors/factory1/line1/pressure
トピック: sensors/factory2/# (ワイルドカード)
コンテンツベース(Content-based):
メッセージの内容(属性値)に基づいてフィルタリングします。AMQPのHeader Exchangeがこれに相当します。
フィルタ条件: temperature > 80 AND location = "factory1"
→ このフィルタに合致するメッセージのみ配信
QoSレベル(MQTTの例)
| QoS | 配信保証 | 説明 | 用途 |
|---|---|---|---|
| 0 | At most once | 最大1回(ロスあり) | センサーの定期データ |
| 1 | At least once | 最低1回(重複あり) | アラート通知 |
| 2 | Exactly once | 正確に1回 | 課金・取引データ |
トピック設計のベストプラクティス
IoT向けMQTTトピックの設計例:
# 推奨構造: {組織}/{場所}/{デバイスID}/{データ種別}
kobesoft/osaka/sensor_001/temperature
kobesoft/osaka/sensor_001/humidity
kobesoft/osaka/gateway_01/status
# デバイスへのコマンドは逆向きのトピックを使う
kobesoft/osaka/sensor_001/cmd/reboot
kobesoft/osaka/sensor_001/cmd/config
# ワイルドカード購読
kobesoft/osaka/# → 大阪拠点のすべて
kobesoft/+/sensor_001/+ → すべての拠点のsensor_001のすべてのデータ
動作原理
メッセージフローの詳細
Publisher Broker Subscriber
| | |
|-- CONNECT -------------->| |
|<-- CONNACK --------------| |
| |<-- CONNECT ---------------|
| |--- CONNACK -------------->|
| |<-- SUBSCRIBE (temp/#) ----|
| |--- SUBACK --------------->|
| | |
|-- PUBLISH (temp/room1) ->| |
| |--- PUBLISH (temp/room1) ->| トピック一致
|-- PUBLISH (pressure/1) ->| | (配信されない)
| | |
|-- DISCONNECT ----------->| |
保持メッセージ(Retained Message)
Pub/Subの「時間的分離」を補完する機能です。ブローカーはトピックごとに最後のメッセージを保持し、新しい購読者が接続した瞬間に最新値を受け取ることができます。
# Python paho-mqtt で retain メッセージを送信
import paho.mqtt.client as mqtt
client = mqtt.Client()
client.connect("broker.example.com", 1883)
# retain=True でブローカーに保持させる
client.publish(
topic="devices/sensor_01/status",
payload="online",
qos=1,
retain=True # 新規購読者に即座に配信
)
Fan-outパターン
1つのパブリッシュを複数の購読者に同時配信するFan-outは、Pub/Subの重要な特性です。
Publisher: PUBLISH temperature=25.3
↓
Broker
├──→ Subscriber A(ダッシュボード表示)
├──→ Subscriber B(データベース保存)
├──→ Subscriber C(アラート判定)
└──→ Subscriber D(ML推論エンジン)
追加の購読者を加えてもパブリッシャーのコードを変更する必要がない点が、スケーラビリティの核心です。
シャードパターン(スケールアウト)
高スループット環境では、トピックをシャード(分割)して並列処理します。Apache Kafkaのパーティションがこれに相当します。
sensors/factory1/line1/temp → Partition 0 → Worker A
sensors/factory1/line2/temp → Partition 1 → Worker B
sensors/factory2/line1/temp → Partition 2 → Worker C
用途・ユースケース
大規模センサーデータ収集
# AWS IoT Core + Python でのPub/Sub実装
import json
import time
from awscrt import mqtt
from awsiot import mqtt_connection_builder
# AWS IoT Core への接続
connection = mqtt_connection_builder.mtls_from_path(
endpoint="xxxxxxx.iot.ap-northeast-1.amazonaws.com",
cert_filepath="device_cert.pem",
pri_key_filepath="device_key.pem",
ca_filepath="root_ca.pem",
client_id="factory_sensor_001"
)
connection.connect().result()
# センサーデータを定期Publish
def send_sensor_data():
while True:
payload = {
"timestamp": time.time(),
"temperature": read_temp_sensor(),
"device_id": "factory_sensor_001"
}
connection.publish(
topic="factory/osaka/sensors/temperature",
payload=json.dumps(payload),
qos=mqtt.QoS.AT_LEAST_ONCE
)
time.sleep(5)
Google Cloud Pub/Sub によるイベント処理
# Google Cloud Pub/Sub でのメッセージ受信(サブスクライバー)
from google.cloud import pubsub_v1
subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path(
"my-project", "iot-data-sub"
)
def callback(message):
data = json.loads(message.data.decode())
# 処理ロジック
if data['temperature'] > 80:
send_alert(data)
save_to_database(data)
message.ack() # 処理完了確認
# ストリーミング購読(非ブロッキング)
streaming_pull = subscriber.subscribe(
subscription_path, callback=callback
)
デバイスシャドウ / ツインパターン
クラウド上にデバイスの「仮想コピー(ツイン)」を持つパターンでもPub/Subが使われます。
# デバイスが現在の状態を報告
Publish: devices/sensor_01/state/reported
{
"temperature": 25.3,
"firmware": "v1.2.0",
"battery": 87
}
# クラウドが設定変更を送信
Publish: devices/sensor_01/state/desired
{
"report_interval_sec": 300,
"temperature_alert_threshold": 80
}
# デバイスが変更を適用した後、reportedを更新
イベント駆動アーキテクチャ
マイクロサービス間の非同期連携にもPub/Subを活用します。
[認証サービス] → (user.registered) → [メール送信サービス]
→ [プロファイル初期化サービス]
→ [監査ログサービス]
実装・開発のポイント
実装の選択基準
| 要件 | 推奨実装 |
|---|---|
| IoTデバイス(組み込み) | MQTT(mosquitto, EMQX) |
| 高スループット(10万msg/sec超) | Apache Kafka |
| エンタープライズ信頼性 | AMQP(RabbitMQ) |
| クラウドネイティブ(GCP) | Google Cloud Pub/Sub |
| クラウドネイティブ(AWS) | AWS IoT Core / SNS+SQS |
| リアルタイムWeb | WebSocket |
| 産業機器 | OPC UA Pub/Sub |
バックプレッシャーの管理
購読者の処理が追いつかない場合、キューが溢れる問題(バックプレッシャー)が起きます。
# RabbitMQ: prefetch_countでバックプレッシャー制御
channel.basic_qos(prefetch_count=1) # 1件ずつ処理
# Kafka: コンシューマーラグを監視
# kafka-consumer-groups.sh --describe --group my-group
メッセージのべき等性
Pub/Subではメッセージが重複配信される場合があります(at-least-once)。処理側でべき等性を確保します。
import redis
redis_client = redis.Redis()
def idempotent_callback(message):
msg_id = message.message_id
# 処理済みチェック(Redisに60秒キャッシュ)
if redis_client.get(f"processed:{msg_id}"):
message.ack()
return
process_message(message.data)
# 処理済みマーク
redis_client.setex(f"processed:{msg_id}", 60, "1")
message.ack()
メッセージスキーマの管理
複数のサービスが同一トピックを購読する場合、メッセージ形式の変更が問題になります。Schema Registry(Confluent)やProtocol Buffersを使ってスキーマを管理します。
// センサーデータのProtobufスキーマ
syntax = "proto3";
message SensorData {
string device_id = 1;
double temperature = 2;
double humidity = 3;
int64 timestamp_ms = 4;
}
他技術との比較
| 項目 | Pub/Sub(一般) | リクエスト/レスポンス | イベントストリーミング |
|---|---|---|---|
| 代表例 | MQTT、AMQP | HTTP REST、gRPC | Apache Kafka |
| 結合度 | 疎結合 | 密結合 | 疎結合 |
| スケーラビリティ | 高 | 中 | 高(水平スケール) |
| メッセージ保持 | 短期〜中期 | なし | 長期(再生可能) |
| リアルタイム性 | 高 | 高 | 中(バッチ的) |
| 複雑度 | 中 | 低 | 高 |
Pub/SubモデルはMQTT・AMQP・OPC UAなど多くのIoTプロトコルの基盤となる設計パターンです。WebSocketと組み合わせることで、デバイスからブラウザダッシュボードまでシームレスなデータフローを構築できます。