IoTプロトコル

Pub/Sub

発行者と購読者を分離する通信モデル。

概要

Pub/Sub(Publish/Subscribe、パブリッシュ/サブスクライブ)は、情報の発行者(Publisher)と購読者(Subscriber)を直接接続せずに分離するメッセージング通信パターンです。両者はブローカー(仲介者)またはメッセージバスを通じて間接的にやり取りし、互いの存在を知る必要がありません。

この「疎結合(Loose Coupling)」が最大の特徴であり、以下の3つの次元で分離を実現します。

  • 空間的分離:発行者と購読者がお互いのIPアドレスを知らなくてよい
  • 時間的分離:発行者がメッセージを送った時点で購読者がオンラインである必要がない
  • 同期的分離:発行者はメッセージを送った後、応答を待つ必要がない

IoT分野では、数千〜数万台のセンサーデバイスからのデータを複数のバックエンドサービスが処理する構成においてPub/Subが基本アーキテクチャとなっています。MQTTはPub/Subの最も有名な実装の一つであり、他にもAMQPOPC 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配信保証説明用途
0At most once最大1回(ロスあり)センサーの定期データ
1At least once最低1回(重複あり)アラート通知
2Exactly 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
リアルタイムWebWebSocket
産業機器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、AMQPHTTP REST、gRPCApache Kafka
結合度疎結合密結合疎結合
スケーラビリティ高(水平スケール)
メッセージ保持短期〜中期なし長期(再生可能)
リアルタイム性中(バッチ的)
複雑度

Pub/SubモデルはMQTTAMQPOPC UAなど多くのIoTプロトコルの基盤となる設計パターンです。WebSocketと組み合わせることで、デバイスからブラウザダッシュボードまでシームレスなデータフローを構築できます。

関連用語

参考リンク