スマートメーター × Confluent Cloud リアルタイムデモ(東京23区版)修正版
すみません、自分、雰囲気でスマートメーターのデモを作ってました。
ということで有識者の指摘を受けて色々とBobにより修正した版がこの記事です。
修正箇所は、こちらのCHANGELOGにまとめています。
Confluent Cloud(Kafka)を使って、スマートメーターの擬似データをリアルタイムに配信・集計・可視化するデモです。
東京都 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)」
0-3. 「掲示板型」vs「郵便・書留型」で直感的に伝える
| 郵便・書留型(従来型MQ) | 掲示板型(Kafka) | |
|---|---|---|
| 送り先 | 1対1。誰に届けるか事前指定 | 誰でも好きなときに読みに来られる |
| コピー | 複数に届けるにはコピーが必要 | 1件書けば何人でも独立して読める |
| 過去データ | 読んだら消える(通常) | 保持期間内はいつでも再読可能 |
| 大量投稿 | 突発的な大量投稿でキューが詰まる | パーティション追加でスケールアウト |
| 向いている場面 | 送り先・責任範囲が明確な1対1連携 | 多数の部署・システムがデータを共有利用する場面 |
Kafkaを選ぶ価値は「読む人が多い」「投稿ペースが速い」「突発的な大量投稿がある」の3つが揃うときに最大化されます。
スマートメーターはまさにこの3条件を満たす代表的ユースケースです。
目次
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図が詳細すぎてレンダリングできなかった場合に備えて、画像として貼っておきます。↓)

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.py と 4_dashboard_server_tokyo23.py の KAFKA_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 秒後に自動再接続 |
参考リンク
- Confluent Cloud — Kafka マネージドサービス
- Leaflet.js — OSS 地図ライブラリ
- OpenStreetMap Nominatim — 座標データ取得に使用

