SQL Database からの参照データを Azure Stream Analytics ジョブに使用する

リファレンスデータは、ストリーミングデータと結合して、販売イベントのストリームに商品情報を追加するなど、静的またはゆっくりと変化するデータセットです。 Azure Stream Analyticsは参照データのソースとしてAzure SQL Databaseをサポートしているため、リアルタイム入力とデータを検索・組み合わせることができます。

この記事では、AzureポータルとVisual Studioの両方をStream Analyticsツールを使って、Stream Analyticsの仕事の参照データ入力としてAzure SQL Databaseを設定する方法を示します。

Azureポータルを使ってSQL Database参照データを追加

Azureポータルを使って、Azure SQL Databaseをリファレンス入力ソースとして追加するには、以下の手順をご利用ください。

ポータルの前提条件

  1. Stream Analytics ジョブを作成します。

  2. Stream Analyticsジョブ用のストレージアカウントを作成してください。

    重要

    Azure Stream Analyticsはこのストレージアカウント内にスナップショットを保持します。 保持ポリシーを設定する際は、選択した期間にストリームアナリティクスの仕事で望む回復期間が含まれていることを確認してください。

  3. Stream Analyticsの仕事で参照データとして使うデータセットでAzure SQL Databaseを作成しましょう。

SQL Database の参照データ入力を定義する

  1. Stream Analytics ジョブで、 [ジョブ トポロジ][入力] を選択します。 「 参照入力を追加」を選択し、次に 「SQL Database」を選択してください。

    Stream Analytics Inputsのパネルのスクリーンショットで、「Add reference input」を選択しており、ドロップダウンリストに「Blob storage」と「SQL Database」の値が表示されます。

  2. Stream Analyticsの入力設定を入力してください。 データベース名、サーバー名、そしてデータベースのサインイン資格を選択してください。 参照データ入力を定期的に更新するには、「 On 」を選択し、DD:HH:MMでリフレッシュレートを指定します。 リフレッシュレートが短い大規模なデータセットでは、デルタクエリは開始時間、 @deltaStartTime、終了時間 @deltaEndTimeの間に挿入または削除されたSQL Database内のすべての行を取得し、参照データの変更を追跡します。

    詳細については、「 差分クエリ」を参照してください。

    SQLデータベースの新しい入力ページのスクリーンショットで、左側に設定フォーム、右側にスナップショットクエリがあります。

  3. SQL クエリ エディターでスナップショット クエリをテストします。 詳細については、「AzureポータルのSQLクエリエディターを使ってデータを接続・照会する」をご覧ください。

ジョブ設定でストレージアカウントを指定してください

ストレージアカウント設定「設定」から「ストレージアカウントを追加」を選択してください。

右側のパネルにある「ストレージアカウント追加」ボタンがあるストレージアカウント設定のスクリーンショットです。

ジョブを開始

  1. 他の入力、出力、クエリを設定したら、Stream Analyticsジョブを開始します。

Visual Studioを使ってSQLデータベースの参照データを追加

Visual Studioを使ってAzure SQL Databaseをリファレンス入力ソースとして追加するには、以下の手順をご利用ください:

Visual Studio の前提条件

  1. Visual Studio 用の Stream Analytics ツールをインストールします。 ストリーム分析ツールは以下のVisual Studioバージョンをサポートしています:

    • Visual Studio 2015
    • Visual Studio 2019
  2. Visual Studio 用 Stream Analytics ツールのクイック スタートで理解を深めます。

  3. ストレージ アカウントを作成します。

    重要

    Azure Stream Analyticsはこのストレージアカウント内にスナップショットを保持します。 保持ポリシーを設定する際は、選択した期間にストリームアナリティクスの仕事で望む回復期間が含まれていることを確認してください。

SQL Database テーブルの作成

SQL Server Management Studio を使用して、参照データを格納するためのテーブルを作成します。 詳しくは、SSMS を使用した最初の Azure SQL Database の設計に関するチュートリアルをご覧ください。

次の文で例表が作成されます:

create table chemicals(Id Bigint,Name Nvarchar(max),FullName Nvarchar(max));

サブスクリプションの選択

  1. Visual Studio の [表示] メニューで [サーバー エクスプローラー] を選択します。

  2. Azureを選択して長押し(または右クリック)し、「Microsoft Azureサブスクリプションに繋がる」を選択し、Azureアカウントでサインインします。

Stream Analytics プロジェクトを作成する

  1. [ファイル]>[新しいプロジェクト] の順に選択します。

  2. テンプレートリストで「Stream Analytics」を選択し、次に「Azure Stream Analytics Application」を選択します。

  3. プロジェクト ロケーションソリューション名を入力し、 OKを選択してください。

    Stream AnalyticsテンプレートとAzure Stream Analytics Applicationを選択し、Name、Location、Solution名のボックスがハイライトされたNew Projectダイアログのスクリーンショットです。

SQL Database の参照データ入力を定義する

  1. 新しい入力を作成します。

    入力を選択して表示された新しいアイテムを追加するダイアログのスクリーンショット。

  2. ソリューション エクスプローラーInput.jsonを開きます。

  3. [Stream Analytics Input Configuration](Stream Analytics の入力構成) を設定します。 データベース名、サーバー名、リフレッシュタイプ、リフレッシュレートを入力します。 DD:HH:MM の形式で更新間隔を指定します。

    ストリームアナリティクス入力設定のスクリーンショットで、ドロップダウンリストから入力または選択した値が確認されています。

    一度だけ実行」または「定期的に実行」を選択すると、Visual Studio Input.json ファイルノードの下にあるプロジェクトの中に「[Input Alias].snapshot.sql」という名前のSQL CodeBehindファイルが生成されます。

    SQL CodeBehind ファイル 'Chemicals.snapshot.sql' が強調表示された ソリューション エクスプローラー のスクリーンショット。

    Deltaで定期的にリフレッシュを選択すると、Visual Studio 2つのSQLコードビハインドファイル([Input Alias].snapshot.sql[Input Alias].delta.sql を生成します。

    SQL CodeBehindファイルChemicals.delta.sqlとChemicals.snapshot.sql強調表示したソリューション エクスプローラーのスクリーンショットです。

  4. エディターで SQL ファイルを開き、SQL クエリを記述します。

  5. Visual Studio 2019を使っていてSQL Server Data Toolsをインストールしている場合は、「実行」を選択してクエリをテストできます。 ウィザードが開いてSQLデータベースに接続し、下部のウィンドウにクエリ結果が表示されます。

ストレージ アカウントを指定する

JobConfig.json 開いて、SQL参照スナップショットを保存するストレージアカウントを指定します。

Stream Analyticsジョブ設定のスクリーンショットで、デフォルト設定とグローバルストレージ設定がハイライトされています。

ローカル環境でテストして Azure にデプロイする

ジョブをAzureにデプロイする前に、ローカルでクエリロジックをライブ入力データと照らしてテストできます。 この機能の詳細については、「Visual StudioのAzure Stream Analyticsツールを使ってローカルでライブデータをテスト(プレビュー版)」をご覧ください。 テストが終わったら、「Azureに提出する」を選択してください。 仕事の始め方については、「Visual StudioのAzure Stream Analyticsクイックスタートツールを使ってストリーム分析を作成」のジョブをご覧ください。

デルタ クエリ

デルタクエリを使う場合は、Azure SQL Databaseのテンポラルテーブルを使いましょう。

  1. Azure SQL Database にテンポラル テーブルを作成します。

       CREATE TABLE DeviceTemporal
       (
          [DeviceId] int NOT NULL PRIMARY KEY CLUSTERED
          , [GroupDeviceId] nvarchar(100) NOT NULL
          , [Description] nvarchar(100) NOT NULL
          , [ValidFrom] datetime2 (0) GENERATED ALWAYS AS ROW START
          , [ValidTo] datetime2 (0) GENERATED ALWAYS AS ROW END
          , PERIOD FOR SYSTEM_TIME (ValidFrom, ValidTo)
       )
       WITH (SYSTEM_VERSIONING = ON (HISTORY_TABLE = dbo.DeviceHistory));  -- DeviceHistory table will be used in Delta query
    
  2. スナップショット クエリを作成します。

    @snapshotTimeパラメータを使って、Stream Analyticsランタイムにシステム時の有効なSQLデータベースのテンポラルテーブルから参照データセットを取得するよう指示します。 このパラメータを提供しなければ、クロックの歪みにより基準基準データセットが不正確になるリスクがあります。 以下の例は完全なスナップショットクエリを示しています:

       SELECT DeviceId, GroupDeviceId, [Description]
       FROM dbo.DeviceTemporal
       FOR SYSTEM_TIME AS OF @snapshotTime
    
  3. デルタ クエリを作成します。

    このクエリは、開始時間、 @deltaStartTime、終了時間 @deltaEndTime内に挿入または削除されたSQL Database内のすべての行を取得します。 デルタ クエリでは、スナップショット クエリと同じ列および列の操作を返す必要があります。 この列は、 行が@deltaStartTime@deltaEndTimeの間に挿入されるか削除するかを定義します。 結果の行には、レコードが挿入された場合は 1、削除された場合は 2 のフラグが設定されます。 また、差分期間内のすべての更新を確実に取得できるようにするため、クエリでは SQL Server 側で watermark も追加する必要があります。 ウォーターマークなしでデルタクエリを使うと、誤った参照データセットになる可能性があります。

    更新されたレコードの場合、テンポラル テーブルでは挿入と削除の操作をキャプチャすることによってブックキーピングが行われます。 ストリームアナリティクスランタイムは、デルタクエリの結果を前のスナップショットに適用し、参照データを最新の状態に保ちます。 以下の例はデルタクエリを示しています:

       SELECT DeviceId, GroupDeviceId, Description, ValidFrom as _watermark_, 1 as _operation_
       FROM dbo.DeviceTemporal
       WHERE ValidFrom BETWEEN @deltaStartTime AND @deltaEndTime   -- records inserted
       UNION
       SELECT DeviceId, GroupDeviceId, Description, ValidTo as _watermark_, 2 as _operation_
       FROM dbo.DeviceHistory   -- table we created in step 1
       WHERE ValidTo BETWEEN @deltaStartTime AND @deltaEndTime     -- record deleted
    

    Stream Analyticsランタイムは、チェックポイントを保存するためにデルタクエリに加えて、定期的にスナップショットクエリを実行することがあります。

    重要

    参照データデルタクエリを使う際は、時間参照データテーブルに同じ更新を複数回行わないでください。 これが誤った結果をもたらす可能性があります。 参考データが誤った結果を出す原因となる例を挙げます:

     UPDATE myTable SET VALUE=2 WHERE ID = 1;
     UPDATE myTable SET VALUE=2 WHERE ID = 1;
    

    正しい例:

     UPDATE myTable SET VALUE = 2 WHERE ID = 1 and not exists (select * from myTable where ID = 1 and value = 2);
    

    この条件により、重複更新が発生しないことを保証します。

クエリをテストする

クエリがStream Analyticsの仕事で参照データとして使われる期待されるデータセットが返ってくるか確認してください。 クエリをテストするには、ポータルのジョブトポロジーセクションの「Inputs」へ行ってください。 次に、SQLデータベースの参照入力で サンプルデータ を選択します。 サンプルが利用可能になったら、ファイルをダウンロードして返されたデータが期待通りかどうかを確認できます。 開発やテストの反復を最適化するには、Visual StudioのStream Analyticsツールを活用してください。 また、まずAzure SQL Databaseから正しい結果を返すために、他のツールを使っても構いません。その後、Stream Analyticsジョブでそのクエリを使うこともできます。

Visual Studio Code を使用してクエリをテストする

Visual Studio Code に Azure Stream Analytics ToolsSQL Server (mssql) をインストールし、ASA プロジェクトを設定します。 詳細については、「クイック スタート: Visual Studio Code で Azure Stream Analytics ジョブを作成する」と SQL Server (mssql) 拡張機能のチュートリアルを参照してください。

  1. SQL 参照データの入力を構成します。

    ReferenceSQLDatabase.json ファイルが表示されるVisual Studio Codeエディタータブのスクリーンショットです。

  2. SQL Serverのアイコンを選択し、「接続を追加」を選択してください。

    左側のパネルのスクリーンショットで、接続を追加オプションがハイライトされています。

  3. 接続情報を入力します。

    データベースとサーバー情報ボックスがハイライトされた接続フォームのスクリーンショットです。

  4. 参照SQLを選択し、長押し(または右クリック)して「 クエリ実行」を選択します。

    クエリ実行オプションがハイライトされたコンテキストメニューのスクリーンショットです。

  5. 接続を選択します。

    下のリストから接続プロファイルを作成するダイアログボックスのスクリーンショットで、1つのリストエントリがハイライトされています。

  6. クエリ結果の確認と検証を行います。

    クエリ検索結果のスクリーンショットはVisual Studio Codeエディタータブにあります。

FAQ

Azure Stream AnalyticsでSQL参照データ入力を使うことで追加コストが発生しますか?

ストリームアナリティクスの仕事では 、ストリーミングユニットあたりの追加コスト はありません。 ただし、Stream Analytics ジョブには、Azure ストレージ アカウントが関連付けられている必要があります。 Stream Analyticsジョブは、ジョブ開始およびリフレッシュ間隔中にSQLデータベースをクエリし、参照データセットを取得し、そのスナップショットをストレージアカウントに保存します。 これらのスナップショットを保存すると、Azureストレージアカウントの価格ページに記載された追加料金が発生します。

参照データのスナップショットがSQLデータベースからクエリされ、Azure Stream Analyticsのジョブで使われているかどうかはどうやってわかりますか?

Azureポータルの「Metrics」の論理名でフィルタリングされた2つの指標で、SQLデータベースの参照データ入力の健全性を監視できます。

  • InputEvents:この指標はSQLデータベース参照データセットから読み込まれたレコード数を測定します。
  • InputEventBytes:このメトリックでは、Stream Analytics ジョブのメモリに読み込まれた参照データ スナップショットのサイズが測定されます。

これらの指標は、ジョブがSQLデータベースにクエリして参照データセットを取得し、それをメモリに読み込むかどうかを示します。

特別なタイプのAzure SQL Databaseが必要ですか?

Azure Stream Analyticsはあらゆる種類のAzure SQL Databaseに対応します。 ただし、参照データ入力のリフレッシュレートを設定するとクエリ負荷に影響することがあります。 デルタクエリのオプションを使うには、Azure SQL Databaseのtemporal tablesを使いましょう。

なぜAzure Stream AnalyticsはAzure Storageアカウントにスナップショットを保存するのですか?

Stream Analytics は、イベント処理が厳密に 1 回だけ実行されることと、イベントが少なくとも 1 回配信されることを保証します。 一時的な問題が仕事に影響する場合は、状態を復元するために少量のリプレイが必要です。 リプレイを有効にするには、これらのスナップショットをAzure Storageアカウントに保存する必要があります。 チェックポイントリプレイの詳細については、Azure Stream Analyticsジョブの「Checkpoint and replayの概念」をご覧ください。