5
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

スマートメーター × Confluent Cloud リアルタイムデモ(東京23区版)を、IBM Bob を使って構築しました Ver.2

5
Last updated at Posted at 2026-07-30

スマートメーター × Confluent Cloud リアルタイムデモ(東京23区版)修正版

すみません、自分、雰囲気でスマートメーターのデモを作ってました。

ということで有識者の指摘を受けて色々とBobにより修正した版がこの記事です。

修正箇所は、こちらのCHANGELOGにまとめています。


Confluent Cloud(Kafka)を使って、スマートメーターの擬似データをリアルタイムに配信・集計・可視化するデモです。

  • ダッシュボード(通常時)
  • 赤いバブルはそれほどなく、画面右上の「受信レート」は無視できるほど小さい。
    スクリーンショット 2026-07-30 21.32.15.png

東京都 23 区・345 台のメーターデータが 10 秒ごとに流れ、ブラウザ上の地図でリアルタイムに変化します。

⚠️ デモ映えのための設定について
本デモの送信間隔は 10秒に設定していますが、実際のスマートメーターは30分値(現行第1世代)または1分値(次世代)を基本とします。
画面の動きを見やすくするためのデモ専用設定です。


0. このデモを見せる前に(背景の共有)

0-1. なぜスマートメーターが「今」必要になったか(電力自由化の背景)

2016年の電力小売全面自由化により、電力事業は「発電・送配電・小売」の3事業に分離されました。
その結果、次の問いに答えなければ料金精算ができなくなりました。

「誰の家に、どの発電会社の電力が、何kWh届いたか」を 30分単位で把握する

月1回の検針(アナログメーター)では対応不可能です。これがスマートメーター全国展開の制度的な起点です。

【電力自由化前】
  大手電力会社(発電・送配電・小売を一体運営)
       ↓ 月1回の検針で十分
  家庭用アナログメーター

【電力自由化後】
  発電会社A / 発電会社B / 太陽光事業者C … (多数の主体が混在)
       ↓ 30分単位で「誰の電力を誰が使ったか」を特定しないと精算できない
  スマートメーター(30分値を自動収集)
       ↓
  送配電会社が管理するデータ収集基盤(HES → Confluent)

0-2. なぜ Confluent が「スマートメーターには」必然なのか(3条件)

Kafka / Confluent が本当に必要なのは、次の 3条件が同時に揃う場合です。
1条件だけなら、従来型MQやESBで十分対応できます。

条件 スマートメーター(次世代)の実態 従来型MQで対応できるか
① 超大規模スループット 東京電力管内だけで約2,900万台。次世代1分値では 484,000件/秒 ❌ 対応不可(キューが詰まる)
② 予測不能なバースト 停電復旧時に蓄積データが一斉送信。後段システムを保護しなければならない ❌ バースト時にダウンリスク
③ マルチコンシューマー 速報ダッシュボード・課金システム・AI分析・家庭向けサービスが同じデータを独立して消費 ❌ 1対1転送になりデータのコピーが必要

デモで訴求するポイント: このデモは上記3条件をすべて体験できます。
「①スループット(345件/10秒 → 本番484,000件/秒)」「②停電バースト(--outage-recovery)」「③マルチコンシューマー(速報/確定/AI)」

  • ダッシュボード(停電復旧直後でデータがバーストした時)
  • バブルが赤くなり、画面右上の「受信レート」が6700件以上に。
    スクリーンショット 2026-07-30 21.30.56.png

0-3. 「掲示板型」vs「郵便・書留型」で直感的に伝える

郵便・書留型(従来型MQ) 掲示板型(Kafka)
送り先 1対1。誰に届けるか事前指定 誰でも好きなときに読みに来られる
コピー 複数に届けるにはコピーが必要 1件書けば何人でも独立して読める
過去データ 読んだら消える(通常) 保持期間内はいつでも再読可能
大量投稿 突発的な大量投稿でキューが詰まる パーティション追加でスケールアウト
向いている場面 送り先・責任範囲が明確な1対1連携 多数の部署・システムがデータを共有利用する場面

Kafkaを選ぶ価値は「読む人が多い」「投稿ペースが速い」「突発的な大量投稿がある」の3つが揃うときに最大化されます。
スマートメーターはまさにこの3条件を満たす代表的ユースケースです。


目次

  1. このデモを見せる前に(背景の共有)
  2. デモ概要
  3. システムアーキテクチャ
  4. ファイル構成
  5. データ仕様
  6. データリネージュ
  7. セットアップ手順
  8. 起動方法
  9. ダッシュボードの見方

1. デモ概要

何が動くか

項目 内容
データソース 擬似スマートメーター(Python で生成)
対象エリア 東京都 23 区
メーター数 345 台(各区 15 台)
送信間隔 10 秒(デモ映え用。実際は30分値/1分値)
可視化 バブル地図(Leaflet.js)でリアルタイム更新
データモデル cumulative_kwh(積算値)を送信。差分計算は ksqlDB(Confluent上) で実施
停電バースト --outage-recovery フラグで停電復旧シナリオを体験可能

動作イメージ

  • 各区のバブルが電力使用量に比例して大きくなる
  • 異常(警告・アラート)が発生するとバブルの色が変わる
  • 右側のランキングがリアルタイムに並び替わる
  • --outage-recovery 実行中はバースト検知インジケーターが点灯し、スループットが急増する

本番スケールとの対比(Confluent の必然性)

指標 本デモ(東京23区) 東京電力管内(現行) 東京電力管内(次世代1分値)
メーター数 345 台 2,900 万台 同左
送信レート(通常時) 34.5 件/秒 約 16,100 件/秒(30分値) 484,000 件/秒
停電復旧バースト(1時間分) 約 12,420 件 約 17.4 億件 17.4 億件
Confluent の役割 デモ: フローの可視化 高スループット吸収 + マルチコンシューマー スケールアウトで吸収

次世代スマートメーター(経産省 Rev5.1, 2024年度〜順次展開)では 1分値収集が標準仕様となる。


2. システムアーキテクチャ

データフロー

PC(データ生成・送信)
        │
        │  JSON / SASL_SSL / lz4圧縮
        ▼
Confluent Cloud(Kafka)
  Topic: smart_meter_tokyo23
  パーティション: 23(区数と同じ)
        │
        │  Consumer poll()
        ▼
PC(集計サーバー)
  1秒ウィンドウで区ごとに集計
        │
        │  WebSocket ws://127.0.0.1:8765
        ▼
ブラウザ(地図ダッシュボード)
  Leaflet.js でバブル地図をリアルタイム更新

技術スタック

レイヤー 技術
データ生成 Python 3.12(標準ライブラリのみ)
Kafka クライアント confluent-kafka 2.3+
メッセージブローカー Confluent Cloud(AWS us-east-2)
サーバー Python asyncio + WebSocket(websockets 12+)
フロントエンド Vanilla JS + Leaflet.js v1.9.4
地図タイル CartoDB Dark Matter(無料・API キー不要)

Confluent Cloud 環境

項目
Cluster ID lkc-xqmr3zx
Type STANDARD(AWS us-east-2)
Bootstrap pkc-921jm.us-east-2.aws.confluent.cloud:9092
Topic smart_meter_tokyo23
パーティション数 23
保持期間 7 日間

3. ファイル構成

smart-meter-demo/
├── smart_meter_generator_tokyo23.py   # 擬似データ生成ライブラリ(23区・345台)
├── 2_produce_to_confluent_tokyo23.py  # Kafka Producer
├── 4_dashboard_server_tokyo23.py      # WebSocket サーバー(集計・配信)
├── 4_dashboard_map_tokyo23.html       # ブラウザ ダッシュボード
└── requirements.txt                   # confluent-kafka, websockets

4. データ仕様

生データスキーマ(Kafka に流れるメッセージ)

フィールド 説明
meter_id string CHI-0001 区コード3文字 + 連番4桁
district string Chiyoda 区キー(英語)
timestamp string 2026-07-27T08:00:00.000Z 計測時刻(UTC / ISO 8601 ms 精度)
kwh float 4.823 電力使用量(kWh)
voltage float 100.3 電圧(V)。標準 100V
current float 48.23 電流(A)
meter_status string normal normal / warning / alert

メッセージ例:

{
  "meter_id":     "CHI-0001",
  "district":     "Chiyoda",
  "timestamp":    "2026-07-27T08:00:00.000Z",
  "kwh":          4.823,
  "voltage":      100.3,
  "current":      48.23,
  "meter_status": "normal"
}

電力使用量の計算式

kwh = base_kwh
      × hour_multiplier   # 時間帯乗数(深夜 0.27 〜 夕方ピーク 1.25)
      × individual_factor # 個体差(±30%、一様分布)
      × noise             # ランダムノイズ(±5%)
      × spike             # スパイク(2% 確率で 2〜4 倍)

異常判定ロジック

ALERT   : kwh > base × 3.5  または  voltage < 95V  または  voltage > 107V
WARNING : kwh > base × 2.5  または  voltage < 97V  または  voltage > 105V
NORMAL  : 上記以外

区ごとの基準使用量(base_kwh)

base_kwh base_kwh base_kwh
千代田区 5.2 台東区 4.0 豊島区 4.0
中央区 4.8 墨田区 3.6 北区 3.5
港区 5.0 江東区 3.9 荒川区 3.3
新宿区 4.9 品川区 4.2 板橋区 3.6
文京区 3.8 目黒区 3.7 練馬区 3.2
渋谷区 4.7 大田区 4.1 足立区 3.4
中野区 3.4 世田谷区 3.5 葛飾区 3.3
杉並区 3.3 江戸川区 3.2

千代田・港・新宿など都心部ほど base_kwh が高く設定されています(オフィス・商業施設の多さを反映)。

時間帯乗数

時間帯 乗数 特徴
0〜5 時 0.27〜0.35 深夜・最低水準
6〜8 時 0.50〜0.90 朝の立ち上がり
9〜16 時 0.95〜1.05 日中・ほぼ基準値
17〜20 時 1.10〜1.25 夕方ピーク
21〜23 時 0.55〜1.10 夜・徐々に低下

5. データリネージュ

データは以下の順に変換されながら流れます。

5-1. システム全体フロー


5-2. フィールドレベルのデータリネージュ

(5−2図が詳細すぎてレンダリングできなかった場合に備えて、画像として貼っておきます。↓)
スクリーンショット 2026-07-30 21.16.32.png


5-3. 停電バーストシナリオのデータフロー


5-4. フィールドの変化まとめ

フィールド 生成 Kafka Topic ksqlDB delta 集計後 表示
meter_id ✅ 生成 Key として使用 通過 削除
district ✅ 生成 通過 通過 集計キー 区名・バブル
timestamp ✅ 生成 通過 通過 削除
cumulative_kwh ✅ 生成(積算値) 通過 LAG()の入力 削除
kwh ✅ 生成(増分) 通過 LAG() fallback total/avg/max_kwh バブル半径・KPI
delta_kwh ksqlDBで新規生成 集計に使用
voltage ✅ 生成 通過 通過 削除
current ✅ 生成 通過 通過 削除
meter_status ✅ 生成 通過 通過 warning/alert_count バブル色・KPI
value_type ✅ 生成 通過 通過 削除
is_outage_recovery ✅ 生成 通過 通過 outage_recovery_count バーストバナー
sequence_no ✅ 生成 通過 通過 削除
total_kwh 新規生成(Σ) バブル半径・ランキング
outage_recovery_count 新規生成(COUNT) バーストインジケーター
color / label 定数テーブルから付与 区の色・名前

5-5. end-to-end レイテンシ

通常時の合計レイテンシ: 概ね 1〜2 秒
バースト時: Confluent がバーストを蓄積し後段は通常ペース(1秒集計)で消費するため、後段への遅延影響はなし


6. セットアップ手順

前提条件

  • Python 3.12 以上
  • Confluent Cloud アカウント(無料トライアルあり)
  • インターネット接続(地図タイル取得に必要)

インストール

# リポジトリのクローン
git clone <this-repo>
cd smart-meter-demo

# 仮想環境の作成
python3 -m venv .venv
source .venv/bin/activate

# 依存パッケージのインストール
pip install confluent-kafka websockets

Confluent Cloud の準備

# 1. Confluent Cloud でクラスター作成後、API Key を発行
confluent login
confluent kafka cluster list

# 2. トピック作成
confluent kafka topic create smart_meter_tokyo23 \
  --cluster <CLUSTER_ID> \
  --partitions 23 \
  --config retention.ms=604800000

接続情報の設定

2_produce_to_confluent_tokyo23.py4_dashboard_server_tokyo23.pyKAFKA_CONFIG を自分の環境に書き換えます:

KAFKA_CONFIG = {
    "bootstrap.servers": "<Bootstrap Server>",
    "security.protocol": "SASL_SSL",
    "sasl.mechanisms":   "PLAIN",
    "sasl.username":     "<API Key>",
    "sasl.password":     "<API Secret>",
}

7. 起動方法

ターミナル A — データ送信(Producer)

cd smart-meter-demo
.venv/bin/python 2_produce_to_confluent_tokyo23.py
# → 10秒ごとに 345件 を Confluent Kafka へ送信

ターミナル B — WebSocket サーバー

cd smart-meter-demo
.venv/bin/python 4_dashboard_server_tokyo23.py --local-aggregate
# → ws://127.0.0.1:8765 で起動

ブラウザ — ダッシュボードを開く

open smart-meter-demo/4_dashboard_map_tokyo23.html

オフラインモード(Confluent 不要):

.venv/bin/python 4_dashboard_server_tokyo23.py --offline

8. ダッシュボードの見方

┌──────────────────────────────────────────────────────┐
│ タイトル                         最終更新  接続状態  │
├──────┬──────────┬────────┬─────────┬────────────────┤
│総電力│メーター数│ 警告  │ アラート │  最高使用区    │  ← KPI
├──────┴──────────┴────────┴─────────┴──────────┬──────┤
│                                                │ランキ│
│         東京23区 地図(バブル表示)             │  ング│
│                                                │      │
└────────────────────────────────────────────────┴──────┘

バブルの見方

要素 意味
バブルの大きさ total_kwh に比例(大きいほど消費電力が多い)
バブルの色(区固有色) 全メーター正常
バブルの色(黄色) 1 台以上で WARNING 発生中
バブルの色(赤色) 1 台以上で ALERT 発生中

接続状態インジケーター

表示 意味
黄点滅「接続中...」 WebSocket 接続を試みています
緑点灯「接続済み」 正常。10 秒ごとにデータが届いています
赤点灯「切断」 サーバーとの接続が切れました。3 秒後に自動再接続

参考リンク

5
1
2

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
5
1

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?