スマートメーター × Confluent Cloud リアルタイムデモ(東京23区版)
- Confluent Cloud(Kafka)を使って、スマートメーターの擬似データをリアルタイムに配信・集計・可視化するデモです。
- 東京都 23 区・345 台のスマートメーターにて発生したという想定の疑似データが 10 秒ごとに流れ、ブラウザ上の地図でリアルタイムに10秒ごとに変化します。
- 10秒ごとの更新というのは、未来のスマートメーターとして検討されているレベル感としてデモに組み込んだものです。(以下URLご参照)
- https://www.enecho.meti.go.jp/category/electricity_and_gas/electric/summary/regulations/teiatsu_smartmeter_rev5.1_20260327.pdf
参考情報
以下、IBM Community のブログにも同様の内容をPOSTしています。
https://community.ibm.com/community/user/blogs/shumpei-kubo/2026/07/29/confluent-cloud-23ibm-bob
目次
1. デモ概要
何が動くか
| 項目 | 内容 |
|---|---|
| データソース | 擬似スマートメーター(Python で生成) |
| 対象エリア | 東京都 23 区 |
| メーター数 | 345 台(各区 15 台) |
| 送信間隔 | 10 秒ごと |
| 可視化 | バブル地図(Leaflet.js)でリアルタイム更新 |
ダッシュボードの動作イメージ
- 各区のバブルが電力使用量に比例して大きくなる
- 異常(警告・アラート)が発生するとバブルの色が変わる
- 右側のランキングがリアルタイムに並び替わる
Confluent Cloud 内のTopicのデータのイメージ
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 環境
| 項目 | 値 |
|---|---|
| Type | STANDARD(AWS us-east-2) |
| 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. データリネージュ
データは以下の順に変換されながら流れます。
[生成レコード] [Kafka Message] [Confluent Topic]
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ meter_id │──KEY──▶│ KEY:meter_id│────────▶│ KEY:meter_id│
│ district │ │ district │ │ district │
│ timestamp │ │ timestamp │ │ timestamp │
│ kwh │ │ kwh │ │ kwh │
│ voltage │ │ voltage │ │ voltage │
│ current │ │ current │ │ current │
│ meter_status│ │ meter_status│ │ meter_status│
└─────────────┘ └─────────────┘ └─────────────┘
345件/10秒 SASL_SSL/lz4 23 partitions
↓ Consumer poll() + 1秒ウィンドウ集計
[集計オブジェクト] [WebSocket JSON]
┌────────────────────┐ ┌────────────────────┐
│ district ←引継│ │ district │
│ total_kwh ←Σ kwh │ │ total_kwh │
│ avg_kwh ←÷count │────ws://127.0.0.1:8765─▶│ avg_kwh │
│ meter_count ←COUNT│ │ meter_count │
│ warning_count←COUNT│ │ warning_count │
│ alert_count ←COUNT│ │ alert_count │
│ max_kwh ←MAX │ │ max_kwh │
│ color/label ←付与 │ │ color / label │
│ ✗ meter_id (削除)│ └────────────────────┘
│ ✗ timestamp (削除)│ 23区分を一括送信
│ ✗ voltage (削除)│ ~3KB/回
│ ✗ current (削除)│
└────────────────────┘
フィールドの変化まとめ
| フィールド | 生成 | Kafka | 集計後 | 表示 |
|---|---|---|---|---|
meter_id |
生成 | Key として使用 | 削除 | — |
district |
生成 | 通過 | 集計キー | 区名表示 |
timestamp |
生成 | 通過 | 削除 | — |
kwh |
生成 | 通過 |
total/avg/max_kwh に集約 |
バブル半径・KPI |
voltage |
生成 | 通過 | 削除 | — |
current |
生成 | 通過 | 削除 | — |
meter_status |
生成 | 通過 |
warning/alert_count に集約 |
バブル色 |
total_kwh |
— | — | 新規生成(Σ kwh) | バブル半径・ランキング |
warning_count |
— | — | 新規生成(COUNT) | バブル色・KPI |
color / label
|
— | — | 定数テーブルから付与 | 区の色・名前 |
end-to-end レイテンシ
データ生成(0ms)
→ Kafka 送信(~50ms、linger.ms=50)
→ Consumer 受信(即時)
→ 集計ウィンドウ(最大 1000ms)
→ WebSocket broadcast(即時)
→ ブラウザ描画(即時)
合計: 概ね 1〜2 秒
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 — 座標データ取得に使用


