Apache StreamPipes¶
レイヤー: IoTレイヤー
目的¶
Apache StreamPipesはFlexGalaxy.AIのIIoTデータ処理エンジンであり、アカウントスコープのデータ処理、ダッシュボードの可視化、およびアプリケーション統合のためにThingIOによってラップされています。StreamPipesはThingsBoardの組み込みルールエンジンを超える複雑なストリーム分析を処理します — パターン検出、異常分析、機械学習推論、およびマルチソースデータ相関です。
StreamPipesはThingsBoardを補完します。ThingsBoardはデバイス接続と基本的なアラート(DeviceAdmin経由で表示)を担当し、StreamPipesは高度な分析(ThingIO経由で表示)を担当します。
責務の分担¶
懸念事項 |
担当 |
理由 |
|---|---|---|
閾値アラート("バッテリー < 20%") |
ThingsBoardルールエンジン |
シンプル、低レイテンシー、外部依存なし |
パターン検出("バッテリー消耗率の増加") |
StreamPipes |
ウィンドウ計算が必要 |
マルチデバイス相関("ゾーンBでの混雑") |
StreamPipes |
クロスデバイス集約が必要 |
テレメトリに対するMLモデル推論 |
StreamPipes |
MLランタイム統合 |
異常検出 |
StreamPipes |
タイムウィンドウにわたる統計分析 |
Kafkaへのイベントルーティング |
ThingsBoardルールエンジン |
ネイティブ統合、低オーバーヘッド |
ThingsBoardとの統合¶
Devices ──MQTT──► ThingsBoard ──Rule Engine──► Kafka ──► StreamPipes
│ │
│ ▼
│ Analytics results
│ │
▼ ▼
Time-series DB Kafka: sp.analytics.*
(raw storage) │
▼
Platform services
(alerts, dashboards)
データフロー¶
デバイスがMQTT/CoAP/HTTP経由でThingsBoardにテレメトリを送信
ThingsBoardが生テレメトリを時系列データベースに保存
ThingsBoardルールエンジンがKafkaトピック(
tb.telemetry.*)にパブリッシュStreamPipesアダプターがKafkaトピックから消費
StreamPipesパイプラインがデータを処理(フィルタリング、変換、分析)
結果はKafka(
sp.analytics.*)にパブリッシュされ、プラットフォームが消費します
Kafkaトピック¶
方向 |
トピック |
内容 |
|---|---|---|
TB → SP |
|
デバイス位置情報の更新 |
TB → SP |
|
バッテリーレベルの読み取り |
TB → SP |
|
汎用センサーデータ |
TB → SP |
|
タスク進捗イベント |
SP → プラットフォーム |
|
検出された異常 |
SP → プラットフォーム |
|
パターン検出結果 |
SP → プラットフォーム |
|
カバレッジ分析結果 |
SP → プラットフォーム |
|
ゾーン混雑アラート |
パイプラインの例¶
バッテリー劣化検出¶
[Kafka Adapter] [Trend Detector] [Alert Sink]
tb.telemetry.battery ──► Window: 1 hour ──► sp.analytics.anomaly
Detect: drain rate
increasing > 2x
normal
バッテリー消耗率が通常のパターンを超えて加速しているデバイスを検出し、ハードウェアの劣化または過剰なワークロードの可能性を示します。
ゾーン混雑分析¶
[Kafka Adapter] [Geo Aggregation] [Threshold Filter] [Alert Sink]
tb.telemetry.position ──► Group by: zone ──► Count > threshold ──► sp.analytics.congestion
Window: 5 min per zone type
スライディングタイムウィンドウ内でゾーンごとのデバイス数をカウントします。デバイス密度がゾーン容量の閾値を超えると、プランナーが経路変更に使用できる混雑アラートをトリガーします。
清掃カバレッジ検証¶
[Kafka Adapter] [Area Calculator] [Gap Detector] [Result Sink]
tb.telemetry.position ──► Calculate swept ──► Compare to zone ──► sp.analytics.coverage
(cleaning robots only) area per zone total area
Window: per shift Flag uncovered
regions
ClearJanitor向け:ロボットの軌跡から実際の清掃カバレッジを計算し、各ゾーン内の未清掃エリアを特定します。
パイプラインビルダー¶
StreamPipesは、プラットフォームオペレーターがアクセスできるビジュアルパイプラインビルダーを提供します。これにより、コードなしでカスタム分析パイプラインを作成できます:
アダプター — データソース(Kafkaトピック、HTTPエンドポイント、OPC-UA)
処理エレメント — フィルター、変換、集約、ML推論
データシンク — Kafka、データベース、通知、ダッシュボード
StreamPipesには100以上の既製パイプラインエレメントが付属しています。StreamPipes SDKに従ってDockerコンテナとしてカスタムエレメントを追加できます。
デプロイメント¶
StreamPipesは共有EKSクラスター上で稼働します:
コンポーネント |
目的 |
|---|---|
バックエンド |
コアAPI、パイプライン管理、ユーザー管理 |
UI |
Webベースのパイプラインビルダーとモニタリング |
Extensions (all-jvm) |
既製のアダプター、プロセッサー、およびシンク |
Consul |
パイプラインエレメントのサービスディスカバリ |
CouchDB |
パイプラインおよびユーザーデータのストレージ |
InfluxDB |
内部メトリクスおよびパイプライン結果 |
リソース要件¶
規模 |
パイプライン数 |
CPU |
メモリ |
|---|---|---|---|
小規模(開発用) |
< 10 |
2コア |
4 GB |
中規模 |
10–50 |
4コア |
8 GB |
大規模 |
50+ |
8+コア |
16+ GB |
StreamPipes consumes from the shared Strimzi Kafka cluster (KRaft mode, 3
brokers) on the TrustMint EKS cluster. ThingsBoard publishes raw telemetry
to thingsboard.* topics; StreamPipes consumes these and publishes analytics
results to sp.analytics.* topics. ThingIO consumes both topic sets for its
Data API.
責務の境界¶
StreamPipesが担当する範囲:
基本的な閾値を超える複雑なストリーム分析
マルチソースデータの相関と集約
テレメトリデータに対する機械学習モデル推論
カスタム分析用のビジュアルパイプラインビルダー
パターン検出と異常分析
StreamPipesが担当しない範囲:
デバイス接続 → ThingsBoard(DeviceAdmin経由)
生テレメトリの保存 → ThingsBoard
シンプルな閾値アラート → ThingsBoardルールエンジン(DeviceAdmin経由)
OTAアップデート → hawkBit(OTAForge経由)
デバイスプロビジョニング → DeviceAdmin
ダッシュボードのアカウントスコーピング → ThingIO(StreamPipesをアクセス制御でラップ)