メインコンテンツまでスキップ

スタンドアロンのストリーミングテーブルを使用する

スタンドアロンの ストリーミングテーブル は、LakeFlow Pipelinesの外部で定義され、ストリーミングまたは増分データ処理の追加サポートを備えたUnity Catalogに登録されたテーブルです。各ストリーミングテーブルに対してパイプラインが自動的に作成されます。Kafkaおよびクラウドオブジェクトストレージからの増分データロードには、ストリーミングテーブルを使用できます。

スタンドアロン ストリーミング テーブルはDatabricks SQLウェアハウス、またはサーバレス 一般コンピュートで実行されているノートブックから作成および更新できます。 2 つのコンピュート オプションの違いの詳細については、 「スタンドアロン パイプラインの要件」を参照してください。

ノートブックからPythonを使用してスタンドアロンのストリーミングテーブルを作成および更新するには、スタンドアロンパイプラインでPythonを使用するを参照してください。

注記

Delta Lake テーブルをストリーミング ソースとシンクとして使用する方法については、Delta Lake テーブル ストリーミングの読み取りと書き込みを参照してください。

要件​

スタンドアロン ストリーミング テーブルの作成、更新、クエリに関するコンピュート オプション、権限、その他の要件については、 「スタンドアロン パイプラインの要件」を参照してください。

ストリーミングテーブルの作成​

ストリーミングテーブルは、Databricks SQLのSQLクエリによって定義されます。ストリーミングテーブルを作成すると、ソーステーブルに現在存在するデータを使用してストリーミングテーブルが作成されます。 その後、通常はスケジュールに従ってテーブルを更新し、ソース テーブルに追加されたデータをプルしてストリーミングテーブルに追加します。

ストリーミング テーブルを作成すると、テーブルの所有者とみなされます。

既存のテーブルからストリーミング テーブルを作成するには、次の例のようにCREATE STREAMING TABLEステートメントを使用します。

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT product, price FROM STREAM raw_data;

この場合、ストリーミング テーブルsalesは、1 時間ごとに更新されるスケジュールで、 raw_dataテーブルの特定の列から作成されます。 使用するクエリは ストリーミング クエリである必要があります。ストリーミング セマンティクスを使用してソースから読み取るには、 STREAMキーワードを使用します。

更新に使用されるコンピュート​

CREATE OR REFRESH STREAMING TABLEステートメントを使用してストリーミング テーブルを作成すると、初期データの更新と作成がすぐに開始されます。 これらの操作はDatabricks SQLウェアハウス コンピュートを消費しません。 代わりに、ストリーミング テーブルは作成と更新の両方をサーバレス パイプラインに依存します。 専用のサーバレス パイプラインは、ストリーミング テーブルごとにシステムによって自動的に作成され、管理されます。

Auto Loaderでファイルをロードする​

ボリューム内のファイルからストリーミング テーブルを作成するには、 Auto Loader使用します。 クラウド オブジェクト ストレージからのほとんどのデータ取り込みタスクには、Auto Loader を使用します。Auto Loader とパイプラインは、クラウド ストレージに到着するデータが増え続けるにつれて、それを増分的かつべき等的にロードするように設計されています。

Databricks SQLでAuto Loaderを使用するには、read_files 関数を使用します。次の例は、 Auto Loader を使用して大量の JSON ファイルをストリーミングテーブルに読み込む方法を示しています。

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT * FROM STREAM read_files(
"/Volumes/my_catalog/my_schema/my_volume/path/to/data",
format => "json"
);

クラウド ストレージからデータを読み取るには、Auto Loader を使用することもできます。

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
's3://mybucket/analysis/*/*/*.json',
format => "json"
);

Auto Loaderの詳細については、Auto Loaderとはを参照してください。Auto LoaderでのSQL の使用について、例を挙げて詳しく知りたい場合は、オブジェクトストレージからのデータの読み込み を参照してください。

他のソースからのストリーミング取り込み​

Kafka を含む他のソースからの取り込みの例については、 「パイプラインでのデータのロード」を参照してください。

Auto CDCフローを使用して変更データキャプチャ ( CDC ) を適用する​

FLOW AUTO CDC句を使用して、ソースからストリーミング テーブルへのデータキャプチャ ( CDC ) レコードを処理します。 以前は、 MERGE INTOステートメントは Databricks 上で CDC レコードを処理する際によく使用されていました。しかし、 MERGE INTOレコードの順序がずれているために誤った結果を生成する可能性があり、レコードを並べ替えるために複雑なロジックが必要になります。「変更データキャプチャ」と「スナップショット」を参照してください。

AUTO CDC 順不同のレコードを自動的に処理することで、CDCを簡素化します。レコードを識別するためのキー、順序付けのためのシーケンス列、および結果をSCDタイプ1(直接更新)またはSCDタイプ2(履歴追跡)として保存するかどうかを指定します。

次の例では、SCDタイプ1を使用してCDCの変更を適用するストリーミングテーブルを作成します。

SQL
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
SEQUENCE BY sequenceNum
STORED AS SCD TYPE 1;

以下の例では、SCDタイプ2を使用して変更履歴を保持しています。

SQL
CREATE OR REFRESH STREAMING TABLE target
FLOW AUTO CDC
FROM stream(cdc_data.users)
KEYS (userId)
APPLY AS DELETE WHEN operation = "DELETE"
SEQUENCE BY sequenceNum
COLUMNS * EXCEPT (operation, sequenceNum)
STORED AS SCD TYPE 2;

Auto CDCオプションと動作の詳細については、 「AUTO CDC APIs : パイプラインを使用した変更データ キャプチャの簡素化」を参照してください。 完全な構文リファレンスについては、 CREATE STREAMING TABLE を参照してください。

REPLACE WHERE フローを使用して選択的バッチ置換を適用​

テーブル履歴全体を再処理することなく、FLOW REPLACE WHERE句を使用してストリーミングテーブルの対象サブセットを再計算して上書きします。REPLACE WHERE フローは、結合と集計、遅れて到着するデータ、アップストリームの再処理、スキーマ進化、およびバックフィルなどの増分バッチ処理に最適です。

REPLACE WHERE フローの詳細については、要件、述語オーバーライド、および増分更新を含め、スタンドアロン ストリーミング テーブルの REPLACE WHERE フローを参照してください。

REPLACE USINGフローを使用して部分的なスナップショットの置換を適用する​

備考

ベータ版

REPLACE USING フローはベータ版です。

FLOW REPLACE USING 句を使用して、ストリーミングテーブルを部分的なスナップショットのストリームと同期させます。更新のたびに、REPLACE USINGフローは指定されたキー列に一致するすべての行を置き換え、他のすべての行は変更しません。SEQUENCE BY 列は更新を順序付けするため、更新が順不同で到着した場合でも、キーに対する最大のシーケンスが常に優先されます。例えば:

SQL
CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

BY NAME が必要です。これは、位置ではなく名前で列を一致させます。

REPLACE USING は、スタンドアロンのストリーミングテーブルでも Lakeflow pipelines の場合と同様に動作します。仕組み、シーケンス、期待値、制限事項、および例については、REPLACE USING フローによる部分的なスナップショットの置換を参照してください。スタンドアロンのストリーミングテーブルには、以下の違いが適用されます。

  • SQL でフローを定義します。 CREATE OR REFRESH STREAMING TABLE 上でインライン SQL FLOW REPLACE USING 句を使用して、REPLACE USING フローを作成します。スタンドアロンの CREATE FLOW ステートメントは Lakeflow パイプラインの構成要素であり、スタンドアロンのストリーミングテーブルには使用されません。
  • コンピュートは自動的に管理されます。 スタンドアロンのストリーミングテーブルは、システム管理のServerlessパイプライン上で実行され、Databricks Runtime 18.2 以降が必要です。クラシックコンピュートとServerlessコンピュートのどちらかを選択する必要はありません。

新しいデータのみを取り込む​

デフォルトでは、 read_files関数はテーブルの作成中にソース フォルダー内の既存のデータをすべて読み取り、更新ごとに新しく到着するレコードを処理します。

テーブルの作成時にソース フォルダーに既に存在するデータを取り込まないようにするには、 includeExistingFilesオプションをfalseに設定します。つまり、テーブルの作成後にフォルダーに到着したデータのみが処理されます。例えば:

SQL
CREATE OR REFRESH STREAMING TABLE sales
SCHEDULE EVERY 1 hour
AS SELECT *
FROM STREAM read_files(
'/path/to/files',
includeExistingFiles => false
);

ランタイムバージョン​

ストリーミングテーブルは常に最新の Databricks SQL ランタイムバージョンで実行されます。pipelines.channel テーブルプロパティは、以前は preview または current ランタイム チャンネルを選択するために使用されていましたが、現在はサポートされておらず、効果はありません。既存の定義にこのプロパティが含まれている場合、それは安全に無視されるため、削除する必要はありません。

機密データを非表示にする​

ストリーミング テーブルを使用すると、テーブルにアクセスするユーザーから機密データを隠すことができます。 1 つの方法は、機密性の高い列または行を完全に除外するようにクエリを定義することです。あるいは、クエリを実行するユーザーの権限に基づいて、列マスクまたは行フィルターを適用することもできます。たとえば、グループHumanResourcesDeptに属していないユーザーに対してはtax_id列を非表示にすることができます。これを行うには、ストリーミング テーブルの作成時にROW FILTERおよびMASK構文を使用します。 詳細については、 「行フィルターと列マスク」を参照してください。

ストリーミングテーブルの更新​

ストリーミングテーブルは、更新操作を処理するためにサーバレスパイプラインを自動的に作成および使用します。 更新はパイプラインによって管理され、更新はストリーミング テーブルの作成に使用されるDatabricks SQLウェアハウスによって監視されます。 ストリーミングテーブルは、スケジュールに従って実行するパイプラインを使用して更新できます。

更新がスケジュールされている場合でも、いつでも手動更新を呼び出すことができます。更新は、ストリーミング テーブルとともに自動的に作成された同じパイプラインによって処理されます。

ストリーミング テーブルを更新するには:

SQL
REFRESH STREAMING TABLE sales;

DESCRIBE TABLE EXTENDEDで最新の更新のステータスを確認できます。

注記

タイムトラベルクエリを使用する前に、ストリーミング テーブルを更新する必要がある場合があります。

更新をスケジュールする方法については、 「更新のスケジュール」を参照してください。 スケジュールされた更新には更新通知を設定でき、更新のパフォーマンスモードを設定することもできます。

更新の仕組み​

ストリーミング テーブルの更新では、最後の更新後に到着した新しい行のみが評価され、新しいデータのみが追加されます。

各更新では、ストリーミング テーブルの現在の定義を使用して、この新しいデータを処理します。 ストリーミング テーブル定義を変更しても、既存のデータは自動的に再計算されません。 変更が既存のデータと互換性がない場合は (たとえば、データ型の変更など)、次の更新はエラーで失敗します。

次の例は、ストリーミング テーブル定義への変更が更新動作にどのような影響を与えるかを説明しています。

  • フィルターを削除しても、以前にフィルターされた行は再処理されません。
  • 列プロジェクションの変更は、既存データの処理方法には影響しません。
  • 静的スナップショットによる結合は、初期処理時のスナップショット状態を使用します。 更新されたスナップショットと一致する遅れて到着したデータは無視されます。これにより、ディメンションが遅れると、ファクトが削除される可能性があります。
  • 既存の列の CAST を変更するとエラーが発生します。

既存のストリーミング テーブルでサポートできない方法でデータが変更された場合は、完全な更新を実行できます。

ストリーミングテーブルのフルリフレッシュ​

完全更新では、ソースで利用可能なすべてのデータが最新の定義で再処理されます。完全な 更新 によって既存のデータが切り捨てられるため、データの履歴全体が保持されない、または保持期間が短い ソース ( Kafkaなど) で完全な 更新 を呼び出すことは推奨されません。 ソース内でデータが利用できなくなった場合、古いデータを回復できない可能性があります。

例えば:

SQL
REFRESH STREAMING TABLE sales FULL;

更新をスケジュールおよび監視​

ストリーミングテーブルは、スケジュールに基づいて、またはアップストリームのデータが変更されたときに自動的に更新できます。また、更新のタイムアウト、通知、およびパフォーマンスモードを設定できます。「更新のスケジュール」を参照してください。

ストリーミングテーブルへのアクセスを制御する​

ストリーミング テーブルは、潜在的なプライベート データの公開を回避しながら、データ共有をサポートするための豊富なアクセス制御をサポートしています。 ストリーミング テーブルの所有者またはMANAGE権限を持つユーザーは、他のユーザーにSELECT権限を付与できます。 ストリーミング テーブルへのSELECTアクセス権を持つユーザーは、ストリーミング テーブルによって参照されるテーブルへのSELECTアクセス権を必要としません。 このアクセス制御により、基盤となるデータへのアクセスを制御しながらデータ共有が可能になります。

ストリーミング テーブルの所有者を変更することもできます。

ストリーミングテーブルに権限を付与する​

ストリーミング テーブルへのアクセスを許可するには、 GRANTステートメントを使用します。

SQL
GRANT <privilege_type> ON <st_name> TO <principal>;

privilege_typeは次のいずれかになります。

  • SELECT - ユーザーはストリーミング テーブルをSELECTできます。
  • REFRESH - ユーザーはストリーミング テーブルをREFRESHできます。 更新は所有者の権限を使用して実行されます。

次の例では、ストリーミング テーブルを作成し、ユーザーに選択権限と更新権限を付与します。

SQL
CREATE OR REFRESH STREAMING TABLE st_name AS SELECT * FROM STREAM source_table;

-- Grant read-only access:
GRANT SELECT ON st_name TO read_only_user;

-- Grant read and refresh access:
GRANT SELECT ON st_name TO refresh_user;
GRANT REFRESH ON st_name TO refresh_user;

Unity Catalogのセキュリティ保護可能なオブジェクトに対する権限付与に関する詳細については、 Unity Catalog権限リファレンスを参照してください。

ストリーミングテーブルから権限を取り消す​

ストリーミング テーブルからのアクセスを取り消すには、 REVOKEステートメントを使用します。

SQL
REVOKE privilege_type ON <st_name> FROM principal;

ソース テーブルのSELECT権限が、ストリーミング テーブルの所有者、またはストリーミング テーブルでMANAGEまたはSELECT権限を付与されている他のユーザーから取り消された場合、またはソース テーブルが削除された場合でも、ストリーミング テーブルの所有者またはアクセスを許可されたユーザーは引き続きストリーミング テーブルをクエリできます。 ただし、次の動作が発生します。

  • ストリーミング テーブルの所有者またはストリーミング テーブルにアクセスできなくなった他の人は、そのストリーミング テーブルをREFRESHできなくなり、ストリーミング テーブルは時間の経過とともに古くなります。
  • スケジュールを使用して自動化されている場合、次にスケジュールされているREFRESH失敗するか、実行されません。

次の例では、 read_only_userからSELECT権限を取り消します。

SQL
REVOKE SELECT ON st_name FROM read_only_user;

ストリーミングテーブルの所有者を変更する​

MANAGE 権限を持つユーザーは、スタンドアロンのストリーミングテーブルで、カタログエクスプローラーを使用して新しいオーナーを設定できます。新しい所有者は、自身または サービスプリンシパルユーザー ロールを持つサービスプリンシパルにすることができます。

  1. Databricksワークスペースから、データアイコン。 カタログ をクリックしてカタログ エクスプローラーを開きます。

  2. 更新するストリーミング テーブルを選択します。

  3. 右側のサイドバーの 「このストリーミングテーブルについて」 の下で、 所有者 を見つけてクリックします。鉛筆アイコン。編集。

注記

パイプライン設定で Run as ユーザーを変更して所有者を更新するように指示するメッセージが表示された場合は、ストリーミングテーブルはスタンドアロンテーブルではなく、LakeFlow Pipelinesで定義されています。メッセージには、パイプライン設定へのLinkが含まれており、そこで ラン アズ ユーザーを変更できます。

  1. ストリーミング テーブルの新しい所有者を選択します。

    所有者は、自分が所有するストリーミング テーブルに対するMANAGE権限とSELECT権限を自動的に持ちます。 サービスプリンシパルを自分が所有するストリーミング テーブルの所有者として設定していて、ストリーミング テーブルに対するSELECTまたはMANAGE権限を明示的に持っていない場合、この変更によりストリーミング テーブルへのすべてのアクセスが失われます。 この場合、それらの権限を明示的に付与するように求められます。

    「保存」 時に付与するには、 「MANAGE」権限 と 「SELECT」 権限の両方を選択します。

  2. 所有者を変更するには、 「保存」 をクリックします。

ストリーミングテーブルの所有者が更新されます。 今後のすべての更新は、新しい所有者の ID を使用して実行されます。

所有者がソーステーブルに対する権限を失った場合​

所有者を変更し、新しい所有者がソース テーブルにアクセスできない場合 (または、基礎となるソース テーブルに対するSELECT権限が取り消された場合)、ユーザーは引き続きストリーミング テーブルにクエリを実行できます。 しかし:

  • ストリーミング テーブルをREFRESHすることはできません。
  • 次にスケジュールされているストリーミング テーブルの更新は失敗します。

ソース データにアクセスできなくなると更新ができなくなりますが、既存のストリーミング テーブルの読み取りが直ちに無効になるわけではありません。

ストリーミングテーブルからレコードを完全に削除する​

備考

プレビュー

ストリーミング テーブルでのREORGステートメントのサポートはパブリック プレビュー段階です。

注記
  • ストリーミング テーブルでREORGステートメントを使用するには、 Databricks Runtime 15.4 以降が必要です。
  • REORGステートメントはどのストリーミング テーブルでも使用できますが、削除が有効になっているストリーミング テーブルからレコードを削除する場合にのみ必要です。 コマンドは、投下が有効になっていないストリーミング テーブルで使用した場合には効果がありません。

GDPRコンプライアンスなどの削除を有効にしたストリーミング テーブルの基盤となるストレージからレコードを物理的に削除するには、ストリーミング テーブルのデータに対してvacuum操作を確実に実行するための追加のステップを実行する必要があります。

基礎となるストレージからレコードを物理的に削除するには:

  1. ストリーミング テーブルのレコードを更新または削除します。
  2. APPLY (PURGE)パラメーターを指定して、ストリーミング テーブルに対してREORGステートメントを実行します。 たとえばREORG TABLE <streaming-table-name> APPLY (PURGE); 。
  3. ストリーミングテーブルのデータ保持期間が経過するまで待ちます。 デフォルトのデータ保持期間は 7 日間ですが、 delta.deletedFileRetentionDurationテーブル プロパティを使用して構成できます。「タイムトラベルクエリのデータ保持を構成する」を参照してください。
  4. REFRESH ストリーミングテーブル。「ストリーミング テーブルの更新」を参照してください。 REFRESH操作から 24 時間以内に、レコードが完全に削除されるようにするために必要なVACUUM操作を含むパイプライン メンテナンス タスクが自動的に実行されます。

クエリ履歴を使用して実行を監視する​

クエリ履歴ページを使用すると、クエリの詳細とクエリ プロファイルにアクセスできます。これらは、パフォーマンスの悪いクエリや、ストリーミング テーブルの更新を実行するために使用されるパイプラインのボトルネックを特定するのに役立ちます。 クエリ履歴とクエリ プロファイルで利用できる情報の種類の概要については、 「クエリ履歴」と「クエリ プロファイル」を参照してください。

備考

プレビュー

この機能はパブリック プレビュー段階です。ワークスペース管理者は、 プレビュー ページからこの機能へのアクセスを制御できます。「Databricks プレビューの管理」を参照してください。

ストリーミング テーブルに関連するすべてのステートメントはクエリ履歴に表示されます。 ステートメント ドロップダウン フィルターを使用して、任意のコマンドを選択し、関連するクエリを検査できます。すべてのCREATEステートメントの後には、パイプラインで非同期に実行されるREFRESHステートメントが続きます。REFRESHステートメントには通常、パフォーマンスの最適化に関する情報を提供する詳細なクエリ プランが含まれます。

クエリ履歴 UI でREFRESHステートメントにアクセスするには、次のステップを使用します。

  1. クリック履歴アイコン。左側のサイドバーにある をクリックして、 書き込みー履歴 UIを開きます。
  2. ステートメント ドロップダウン・フィルターから REFRESH チェック・ボックスを選択します。
  3. クエリ ステートメントの名前をクリックすると、クエリの実行時間や集計されたメトリックなどの概要の詳細が表示されます。
  4. クエリ プロファイルを開くには、[クエリ プロファイルを表示 ] をクリックします。クエリ プロファイルのナビゲートの詳細については、「クエリ プロファイル」を参照してください。
  5. 必要に応じて、 [クエリ ソース] セクションのリンクを使用して、関連するクエリまたはパイプラインを開くことができます。

SQLエディターのリンクを使用するか、 SQLウェアハウスに接続されているノートブックからクエリの詳細にアクセスすることもできます。

外部クライアントからストリーミングテーブルにアクセスする​

オープンAPIsをサポートしていない外部のDelta LakeまたはIcebergクライアントからストリーミング テーブルにアクセスするには、互換Modeを使用できます。 互換Modeは、 Delta LakeまたはIcebergクライアントからアクセスできる読み取り専用バージョンのストリーミング テーブルが作成されます。

その他のリソース​