Azure Stream Analytics でのカスタム BLOB 出力のパーティション分割

Azure Stream Analytics は、カスタムのフィールドまたは属性とカスタムの DateTime パス パターンを使用したカスタムの BLOB 出力のパーティション分割をサポートしています。

カスタム フィールドまたは属性

カスタム フィールドまたは入力属性では、出力をよりきめ細かく制御できるようになるため、ダウンストリームのデータ処理およびレポート ワークフローが機能強化されます。

パーティション キーのオプション

入力データのパーティション分割に使用されるパーティション キー (列名) には、 BLOB 名に使用できる任意の文字が含まれる場合があります。 エイリアスと共に使用しない限り、入れ子になったフィールドをパーティション キーとして使用することはできません。 ただし、特定の文字を使用してファイルの階層を作成できます。 たとえば、列を作成し、この列で他の 2 つの列のデータを結合して、一意のパーティション キーを作成するには、次のクエリを使用できます。

SELECT name, id, CONCAT(name, "/", id) AS nameid

パーティション キーは、NVARCHAR(MAX)BIGINTFLOAT、または BIT (互換性レベル 1.2 以上) である必要があります。 DateTimeArrayRecords 型はサポートされていませんが、文字列に変換すると、パーティション キーとして使用できます。 詳細については、Azure Stream Analytics のデータ型に関するページを参照してください。

カスタム フィールドで blob 出力をパーティション分割する

あるジョブが、外部のビデオ ゲーム サービスに接続されたライブ ユーザー セッションから入力データを取得するとします。ここで、取り込まれたデータには、セッションを識別するための列 client_id が含まれます。 データを client_id でパーティション分割するには、ジョブの作成時に BLOB 出力プロパティの パス パターン フィールドにパーティション トークン {client_id} を含めるように設定します。 さまざまな client_id 値を含むデータは Stream Analytics ジョブを通して流れるため、出力データは、フォルダーごとの 1 つの client_id 値に基づいて個別のフォルダーに保存されます。

クライアント ID を含むパス パターンを示すスクリーンショット。

同様に、ジョブ入力が、各センサーに sensor_id が割り当てられた数百万のセンサーからのセンサー データであった場合は、各センサー データを異なるフォルダーにパーティション分割するために、パス パターンは {sensor_id} になります。

REST API を使用すると、その要求に使用される JSON ファイルの出力セクションは次の画像のようになります。

BLOB ストレージの REST API JSON 出力構成を示すスクリーンショット。

ジョブが実行を開始すると、clients コンテナーは次の画像のようになります。

クライアント コンテナーを示すスクリーンショット。

各フォルダーには、各 BLOB に 1 つ以上のレコードが含まれる複数の BLOB を含めることができます。 前の例では、"06000000" のラベルが付けられたフォルダー内に、次のコンテンツを含む 1 つの BLOB が存在します。

クライアント ID レコードを含む BLOB の内容を示すスクリーンショット。

出力パス内の出力をパーティション分割するために使用される列は client_id であったため、BLOB 内の各レコードにはフォルダー名に一致する client_id 列が含まれることに注意してください。

カスタム パーティション キーの制限事項

  1. パス パターンの BLOB 出力プロパティで指定できるカスタムのパーティション キーは 1 つだけです。 次のすべてのパス パターンが有効です。

    • cluster1/{date}/{aFieldInMyData}
    • cluster1/{time}/{aFieldInMyData}
    • cluster1/{aFieldInMyData}
    • cluster1/{date}/{time}/{aFieldInMyData}
  2. 複数の入力フィールドを使用する場合は、 CONCATを使用して BLOB 出力のカスタム パス パーティションのクエリに複合キーを作成できます。 たとえば select concat (col1, col2) as compositeColumn into blobOutput from input です。 その後、Azure Blob Storageのカスタム パスとしてcompositeColumnを指定できます。

  3. パーティション キーは大文字と小文字が区別されないため、Johnjohn などのパーティション キーは同等です。 また、式をパーティション キーとして使用することはできません。 たとえば、{columnA + columnB} では機能しません。

  4. 入力ストリームがパーティション キーのカーディナリティが 8,000 以下のレコードで構成されている場合、Stream Analytics は既存の BLOB にレコードを追加し、必要に応じて新しい BLOB のみを作成します。 カーディナリティが 8,000 を超える場合、Stream Analytics が既存の BLOB に書き込む保証はありません。 Stream Analytics では、同じパーティション キーを持つ任意の数のレコードの新しい BLOB は作成されません。

  5. BLOB 出力が不変として構成されている場合、データが送信されるたびに Stream Analytics によって新しい BLOB が作成されます。

カスタム日付時刻パスパターン

カスタムの DateTime パス パターンを使用すると Hive Streaming 規則に合致する出力形式を指定できます。これにより、Stream Analytics は、ダウンストリーム処理するためのデータを Azure HDInsight と Azure Databricks に送信できるようになります。 カスタムの DateTime パス パターンを簡単に実装するには、BLOB 出力の datetime フィールドで キーワードと書式指定子を使用します。 たとえば {datetime:yyyy} です。

サポートされているトークン

次の書式指定子トークンを単独でまたは組み合わせて使用して、カスタムの DateTime 形式を指定できます。

書式指定子 説明 サンプル時間 2018-01-02T10:06:08 に対する結果
{datetime:yyyy} 4 桁の数値としての年 2018
{datetime:MM} 月番号(01 ~ 12) 01
{datetime:M} 月の範囲:1 から 12 まで 1
{datetime:dd} 日 (01 ~ 31) 02
{datetime:d} 1日から31日までの日 2
{datetime:HH} 24時間表記での時間 (00 ~ 23) 10
{datetime:mm} 分 (00 ~ 60) 06
{datetime:m} 0分から60分まで 6
{datetime:ss} 秒数は00から60まで 08

カスタムの DateTime パターンを使用しない場合は、{date}{time} のトークンを [パス プレフィックス] フィールドに追加して、組み込みの DateTime 形式のドロップダウンを生成できます。

Stream Analytics の古い DateTime 形式を示すスクリーンショット。

DateTime トークンの機能拡張と制限

プレフィックス文字数の制限に達するまで、任意の数のトークン {datetime:<specifier>} をパス パターン内で使用できます。 日時ドロップダウンに既に一覧されている組み合わせ以外の書式指定子を 1 つのトークンの中で組み合わせて使用することはできません。

logs/MM/dd のパス パーティションの場合:

有効な式 無効な式
logs/{datetime:MM}/{datetime:dd} logs/{datetime:MM/dd}

パス プレフィックス内で、同じ書式指定子を複数回使用できます。 トークンは毎回繰り返す必要があります。

Hive Streaming 規則

Hive ストリーミング規則を使用してBlob Storageにカスタム パス パターンを使用できます。この規則では、フォルダー名にcolumn=でフォルダーのラベルが付けられます。

たとえば year={datetime:yyyy}/month={datetime:MM}/day={datetime:dd}/hour={datetime:HH} です。

カスタム出力は Stream Analytics と Hive 間でデータを移植するためのテーブルの変更と手動でのパーティションの追加という煩わしい操作を排除します。 代わりに、以下を使用して、多数のフォルダーを自動的に追加できます。

MSCK REPAIR TABLE while hive.exec.dynamic.partition true

Hive と互換性のある DateTime フォルダー構造を作成する

Stream Analytics Azure portal クイックスタートに従って、ストレージ アカウント、リソース グループ、Stream Analytics ジョブ、入力ソースを作成します。 クイックスタートで使用したものと同じサンプル データを使用します。 サンプル データは GitHub でも入手できます。

次の構成の BLOB 出力シンクを作成します。

Stream Analytics の BLOB 出力シンクの作成を示すスクリーンショット。

完全なパス パターンは次のとおりです。

year={datetime:yyyy}/month={datetime:MM}/day={datetime:dd}

ジョブを開始すると、パス パターンに基づくフォルダー構造が BLOB コンテナー内に作成されます。 日レベルまでドリルダウンできます。

カスタムのパス パターンを含む Stream Analytics の BLOB 出力を示すスクリーンショット。