# RedshiftへのMQTTデータ取り込み

[Amazon Redshift](https://aws.amazon.com/redshift/?nc1=h_ls) は、ペタバイト規模のクラウドデータウェアハウスであり、高性能な分析を目的としたフルマネージドサービスです。PostgreSQLをベースにし、オンライン分析処理（OLAP）に最適化されているため、複雑なクエリの実行や大規模なデータ分析を高速に行えます。EMQXはAmazon Redshiftと直接連携し、IoTデバイスからのMQTTテレメトリをほぼリアルタイムで取り込み、保存することが可能です。

本ページでは、EMQXとRedshift間のデータ統合について包括的に解説し、実際の作成および検証手順を紹介します。

## 動作概要

EMQXのRedshiftデータ統合は組み込み機能であり、MQTTベースのIoTデータストリームをAmazon Redshiftの分散型PostgreSQL互換データウェアハウスに直接取り込みます。EMQXの組み込み[ルールエンジン](./rules.md)を利用することで、複雑なカスタムコードを書かずにIoTデータをRedshiftにストリーミングし、大規模な分析処理が可能です。

以下の図は、EMQXとRedshift間の典型的なデータ統合アーキテクチャを示しています。

![EMQX Integration Redshift](./assets/redshift_architecture.png)

RedshiftへのMQTTデータ取り込みの流れは以下の通りです：

1. **IoTデバイスがEMQXに接続**：IoTデバイスがMQTTプロトコルを介して正常に接続されると、オンラインイベントがトリガーされます。イベントにはデバイスID、送信元IPアドレスなどの情報が含まれます。
2. **メッセージのパブリッシュと受信**：デバイスは特定のトピックにテレメトリやステータスデータをパブリッシュします。EMQXがこれらのメッセージを受信すると、ルールエンジン内でマッチング処理が開始されます。
3. **ルールエンジンによるメッセージ処理**：EMQXのルールエンジンは、トピックやメッセージ内容に基づいて定義されたルールにマッチさせてイベントやメッセージを処理します。処理内容にはデータ変換（例：JSONからSQL用フォーマットへの変換）、フィルタリング、コンテキスト情報によるデータ強化などが含まれ、データベース挿入前に行われます。
4. **Redshiftへの書き込み**：マッチしたルールはSQLベースの取り込みをRedshiftに対してトリガーします。SQLテンプレートを用いて、処理済みデータのフィールドをRedshiftのテーブルおよびカラムにマッピングします。高スループットの取り込みには、Amazon S3からのCOPYコマンドやRedshift Streaming Ingestionを活用し、カラムナーストアに効率的にロードします。RedshiftのクエリオプティマイザとMPP（Massively Parallel Processing）実行エンジンにより、データは即座に分析クエリに利用可能となります。

イベントおよびメッセージデータがRedshiftに書き込まれた後は、以下のような活用が可能です：

- Amazon QuickSight、Grafana、TableauなどのツールとRedshiftを接続し、IoTメトリクスやトレンドを追跡するダッシュボードを構築。
- RedshiftデータをAWSの分析およびAI/MLサービス（例：Amazon SageMaker）と連携し、異常検知やデバイス挙動の予測を実施。
- Redshiftの並列クエリ実行により、大規模なIoTデータセットに対して集計、結合、時系列分析を実行し、過去データとほぼリアルタイムのインサイトを提供。

## 特長と利点

Redshiftとのデータ統合により、以下の特長とメリットをビジネスにもたらします：

- **柔軟なイベント処理**：EMQXのルールエンジンを活用し、Redshiftはデバイスのライフサイクルイベント（接続、切断、ステータス変化）を低レイテンシで保存・処理可能です。RedshiftのMPPクエリエンジンと組み合わせることで、障害検知、異常検知、長期利用傾向の迅速な集計・分析が行えます。
- **メッセージ変換**：メッセージはEMQXルールで広範に処理・変換されてからRedshiftに書き込まれるため、保存データは分析に最適化された状態となります。これによりクエリの複雑さが軽減され、下流処理が効率化されます。
- **SQLテンプレートによる柔軟なデータ操作**：EMQXのSQLテンプレートマッピングを通じて、構造化されたIoTデータをRedshiftのテーブル・カラムに挿入可能です。RedshiftはPostgreSQL互換SQL、JSON用のSUPER型などの半構造化データ型、クエリ最適化のための高度なインデックスをサポートします。カラムナーストレージ、データ圧縮、ゾーンマップにより、大規模データセットのスキャン時間を短縮しクエリを高速化します。
- **ビジネスプロセスの統合**：RedshiftはAWSエコシステムとシームレスに統合されており、IoTデータをAmazon QuickSightなどのBIツール、AWS GlueやAWS Data Pipelineなどの分析サービス、Amazon SageMakerなどのAI/MLサービスに接続可能です。
- **高度な地理空間機能**：RedshiftはGEOMETRYおよびGEOGRAPHY型を通じて地理空間データ型と関数をサポートし、ジオフェンシング、位置情報分析、ルート最適化を実現します。EMQXのリアルタイム取り込みと組み合わせることで、資産追跡、車両監視、位置ベースのイベントトリガーをほぼリアルタイムに行えます。
- **組み込みのメトリクスと監視**：EMQXは各Redshiftシンクのランタイムメトリクスを提供し、RedshiftはAmazon CloudWatchと連携してクラスターのパフォーマンス、クエリ実行メトリクス、ストレージ使用状況を監視可能です。これにより、取り込みから分析までのエンドツーエンドの可観測性を確保します。

## はじめる前に

このセクションでは、Redshift統合の作成を開始する前に必要な準備について説明します。Redshiftクラスターの作成、データベースおよびデータテーブルの作成方法を含みます。

### 前提条件

- EMQXデータ統合の[ルール](./rules.md)に関する知識
- [データ統合](./data-bridges.md)に関する知識

### Amazon Redshiftでのデータベースおよびテーブル作成

EMQXでRedshiftコネクターを設定する前に、Amazon Redshiftクラスター（またはServerlessワークグループ）が稼働していること、そしてIoTデータを格納するスキーマが準備されていることを確認してください。

1. Redshiftクラスターまたはワークグループをデプロイします。[Amazon Redshiftクラスター作成ガイド](https://docs.aws.amazon.com/redshift/latest/mgmt/create-cluster.html)に従い環境を起動してください。

2. データベースユーザーの認証情報を設定します。初期クラスター作成時に、管理者ユーザー（通常は`adminuser`）の資格情報を指定します。

   もしくは、Redshift SQLを使ってEMQX専用のデータベースユーザーを作成します。このユーザーには接続、テーブル作成、読み書きの権限が必要です。例：

   ```sql
   CREATE USER emqx_user PASSWORD 'YourStrongPassword1';
   ```

   詳細は[Redshift入門ガイド](https://docs.aws.amazon.com/redshift/latest/gsg/t_adding_redshift_user_cmd.html)および[ユーザーガイド](https://docs.aws.amazon.com/redshift/latest/dg/r_Users.html)を参照してください。

   後でEMQXのコネクター設定に使用するため、ユーザー名（`emqx_user`）とパスワードを控えておいてください。

3. 任意のPostgreSQL互換クライアント（`psql`、SQL Workbench/J、DBeaverなど）を使用し、ホスト名、ポート、既存のデータベース名（例：デフォルトの`dev`）、ユーザー名、パスワードで[Redshiftエンドポイントに接続](https://docs.aws.amazon.com/redshift/latest/mgmt/cluster-syntax.html)します。

4. 接続後、EMQXからのIoTデータ受け入れ先となる`emqx_data`データベースを作成します。

   ```sql
   CREATE DATABASE emqx_data;
   ```

5. `emqx_data`データベースに接続し、MQTTメッセージおよびクライアントイベントデータを格納するための2つのテーブルを作成します。

   - クライアントID、トピック、ペイロード、作成時刻を保存するデータテーブル`t_mqtt_msg`を以下のSQLで作成します：

     ```sql
     CREATE TABLE t_mqtt_msg (
       id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
       msgid VARCHAR(64),
       sender VARCHAR(64),
       topic  VARCHAR(255),
       qos    INTEGER,
       retain INTEGER,
       -- ペイロードがJSONの場合はSUPER型を検討、そうでなければ大きめのVARCHARを使用
       payload SUPER,
       arrived TIMESTAMPTZ
     );
     ```

   - クライアントのオンライン/オフラインイベントをタイムスタンプ付きで保存する`emqx_client_events`テーブルを以下のSQLで作成します：

     ```sql
     CREATE TABLE emqx_client_events (
       id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
       clientid  VARCHAR(255),
       event     VARCHAR(255),
       created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP
     );
     ```

## Redshiftコネクターの作成

Redshiftシンクを追加する前に、EMQXでRedshiftコネクターを作成する必要があります。コネクターはEMQXがAmazon RedshiftクラスターまたはServerlessワークグループに接続する方法を定義します。

::: tip 注意

Amazon Redshift Serverlessを使用している場合、コネクターが作成され接続が確立されると、データが書き込まれていなくても課金が発生する可能性があります。未使用のコネクターは削除するか、リソースを一時停止して予期せぬコストを回避してください。

:::

1. EMQXダッシュボードで、**Integration** -> **Connector** に移動します。
2. ページ右上の **Create** をクリックします。
3. **Create Connector** ページで **Redshift** を選択し、**Next** をクリックします。
4. コネクター名を入力します。名前は英数字で始まり、英数字、ハイフン、アンダースコアを含めることができます。例：`my_redshift`。
5. Redshift接続情報を入力します：

   - **Server Host**：Redshiftエンドポイントのホスト名（例：`redshift-cluster-1.abc123xyz.us-east-1.redshift.amazonaws.com`）。AWS Redshiftコンソールの**Clusters**または**Workgroups**ページで確認可能です。
   - **Database Name**：EMQXデータを格納する対象データベース。例：`emqx_data`。
   - **Username**：データ挿入権限を持つデータベースユーザー名。例：`emqx_user`。
   - **Password**：`emqx_user`のパスワード。
   - **Enable TLS**：Redshift接続にSSL/TLS暗号化が必要な場合はオンにします（クラウドサービス接続では推奨）。詳細は[外部リソースアクセスのTLS](../../guides/network/overview.md#tls-for-external-resource-access)を参照。
6. 詳細設定（任意）：接続プールサイズ、アイドルタイムアウト、リクエストタイムアウトなどの追加接続プロパティを設定可能です。詳細は[シンクの機能](./data-bridges.md#features-of-sink)を参照してください。
7. **Test Connectivity** をクリックし、EMQXが提供された設定でRedshiftクラスターに正常に接続できるか確認します。
8. **Create** をクリックしてコネクターを保存します。
9. 作成後は以下のいずれかを選択できます：

   - **Back to Connector List** をクリックして全コネクターを表示
   - **Create Rule** をクリックして、このコネクターを使ったルールをすぐに作成し、Redshiftへのデータ転送を設定

   詳細な例は以下を参照してください：

   - [メッセージ保存用Redshiftシンクのルール作成](#create-a-rule-with-redshift-sink-for-message-storage)
   - [イベント記録用Redshiftシンクのルール作成](#create-a-rule-with-redshift-sink-for-events-recording)

## メッセージ保存用Redshiftシンクのルール作成

このセクションでは、ダッシュボードでソースMQTTトピック`t/#`からのメッセージを処理し、処理済みデータを設定済みのRedshiftシンクを介して`t_mqtt_msg`テーブルに保存するルールの作成方法を示します。

1. ダッシュボードの **Integration** -> **Rules** ページに移動します。
2. ページ右上の **Create** をクリックします。
3. ルールIDに `my_rule` を入力し、SQLエディターにルールを入力します。ここではトピック`t/#`のMQTTメッセージをRedshiftに保存するため、ルールのSELECT句でSQLテンプレート内で使用するすべての変数を含むフィールドを選択してください。例：

   ```sql
   SELECT
   *
   FROM
   "t/#"
   ```

   ::: tip

   初心者の方は **SQL Examples** をクリックし、**Enable Test** を有効にしてSQLルールの学習とテストが可能です。

   :::

4. + **Add Action** ボタンをクリックし、ルールでトリガーされるアクションを定義します。このアクションにより、EMQXはルールで処理したデータをRedshiftに送信します。
5. **Type of Action** ドロップダウンからRedshiftを選択し、**Action** ドロップダウンはデフォルトの `Create Action` のままにするか、既存のRedshiftアクションを選択します。この例では新規シンクを作成しルールに追加します。
6. シンクの名前と説明をフォームに入力します。
7. **Connector** ドロップダウンから先ほど作成した `my_redshift` を選択します。新規コネクターを作成する場合はドロップダウン横のボタンをクリックしてください。設定パラメータは[Redshiftコネクターの作成](#create-a-redshift-connector)を参照。
8. **SQL Template** を設定します。以下のSQL文を使用してデータを挿入します。

   注意：これは[プリプロセス済みSQL](./data-bridges.md#prepared-statement)のため、フィールドは引用符で囲まず、文末にセミコロンを付けないでください。

   ```sql
   INSERT INTO t_mqtt_msg (
       msgid,
       topic,
       qos,
       payload,
       arrived
   )
   VALUES (
       ${id},
       ${topic},
       ${qos},
       ${payload},
       timestamp 'epoch' + (${timestamp} :: bigint / 1000) * interval '1 second'
   )
   ```

9. **フォールバックアクション（任意）**：メッセージ配信失敗時の信頼性向上のため、1つ以上のフォールバックアクションを定義可能です。詳細は[フォールバックアクション](./data-bridges.md#fallback-actions)を参照してください。
10. **詳細設定（任意）**：詳細は[シンクの機能](./data-bridges.md#features-of-sink)を参照してください。
11. **Create** をクリックする前に、**Test Connectivity** をクリックしてシンクがRedshiftサーバーに接続可能かテストできます。
12. **Create** ボタンをクリックし、シンクの設定を完了します。新しいシンクが**Action Outputs**に追加されます。
13. **Create Rule** ページで設定内容を確認し、**Save** をクリックしてルールを生成します。

ルール作成後、**Integration** -> **Rules** ページで新規ルールを確認でき、**Action (Sink)** タブで新規Redshiftシンクも確認可能です。

また、**Integration** -> **Flow Designer** でトポロジーを可視化し、トピック`t/#`のメッセージがルール`my_rule`で解析されRedshiftに書き込まれている様子を確認できます。

## イベント記録用Redshiftシンクのルール作成

このセクションでは、クライアントのオンライン/オフライン状態を記録し、イベントデータを設定済みのRedshiftシンクを介して`emqx_client_events`テーブルに保存するルールの作成方法を示します。

手順は[メッセージ保存用Redshiftシンクのルール作成](#create-a-rule-with-redshift-sink-for-message-storage)とほぼ同様ですが、SQLテンプレートとSQLルールが異なります。

オンライン/オフライン状態記録用のSQLルールは以下の通りです：

```sql
SELECT
  *
FROM
  "$events/client_connected", "$events/client_disconnected"
```

イベント記録用SQLテンプレートは以下の通りです：

注意：これは[プリプロセス済みSQL](./data-bridges.md#prepared-statement)のため、フィールドは引用符で囲まず、文末にセミコロンを付けないでください。

```sql
INSERT INTO emqx_client_events(clientid, event, created_at) VALUES (
  ${clientid},
  ${event},
  TO_TIMESTAMP((${timestamp} :: bigint)/1000)
)
```

## ルールのテスト

MQTTXを使用してトピック`t/1`にメッセージを送信し、オンライン/オフラインイベントをトリガーします。

```bash
mqttx pub -i emqx_c -t t/1 -m '{ "msg": "hello Redshift" }'
```

2つのシンクの稼働状況を確認してください。メッセージ保存用シンクでは新規の受信メッセージと送信メッセージが1件ずつあるはずです。イベント記録用シンクでは2件のイベントレコードが確認できます。

`t_mqtt_msg`データテーブルにデータが書き込まれているか確認します。

```bash
emqx_data=# select * from t_mqtt_msg;
 id |              msgid               | sender | topic | qos | retain |            payload
        |       arrived
----+----------------------------------+--------+-------+-----+--------+-------------------------------+---------------------
  1 | 0005F298A0F0AEE2F443000012DC0002 | emqx_c | t/1   |   0 |        | { "msg": "hello Redshift" } | 2023-01-19 07:10:32
(1 row)
```

`emqx_client_events`テーブルにデータが書き込まれているか確認します。

```bash
emqx_data=# select * from emqx_client_events;
 id | clientid |        event        |     created_at
----+----------+---------------------+---------------------
  3 | emqx_c   | client.connected    | 2023-01-19 07:10:32
  4 | emqx_c   | client.disconnected | 2023-01-19 07:10:32
(2 rows)
```
