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)

データフロー

  1. デバイスがMQTT/CoAP/HTTP経由でThingsBoardにテレメトリを送信

  2. ThingsBoardが生テレメトリを時系列データベースに保存

  3. ThingsBoardルールエンジンがKafkaトピック(tb.telemetry.*)にパブリッシュ

  4. StreamPipesアダプターがKafkaトピックから消費

  5. StreamPipesパイプラインがデータを処理(フィルタリング、変換、分析)

  6. 結果はKafka(sp.analytics.*)にパブリッシュされ、プラットフォームが消費します

Kafkaトピック

方向

トピック

内容

TB → SP

tb.telemetry.position

デバイス位置情報の更新

TB → SP

tb.telemetry.battery

バッテリーレベルの読み取り

TB → SP

tb.telemetry.sensor

汎用センサーデータ

TB → SP

tb.telemetry.task

タスク進捗イベント

SP → プラットフォーム

sp.analytics.anomaly

検出された異常

SP → プラットフォーム

sp.analytics.pattern

パターン検出結果

SP → プラットフォーム

sp.analytics.coverage

カバレッジ分析結果

SP → プラットフォーム

sp.analytics.congestion

ゾーン混雑アラート

パイプラインの例

バッテリー劣化検出

[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をアクセス制御でラップ)