ストリームは、Delta テーブルに対する増分クエリでデータのバージョンを自動管理する MaxCompute オブジェクトです。INSERT、UPDATE、DELETE などのデータ操作言語(DML)による変更と、各変更に付随するメタデータを追跡し、増分形式での消費を可能にします。各ストリームはバージョンポインターを保持しており、コンシューマーは常にどの変更が既に処理済みか、またどの変更が新規かを把握できます。
本トピックでは、ストリームの作成、確認、変更、一覧表示、削除、およびクエリ実行に使用する SQL コマンドについて説明します。
仕組み
ストリームは、常に 1 つの Delta テーブルのみに関連付けられます。内部的には、以下の 2 つのバージョンマーカーを維持します。
オフセットバージョン:変更が消費された時点までのデータバージョン。これは、DML 操作内でストリームを読み取った場合にのみ進捗します。
参照テーブルバージョン:関連付けられた Delta テーブルの最新のデータバージョン。テーブルが変更されるたびに自動的に更新されます。
ストリームをクエリするたびに、MaxCompute は半開区間 (オフセットバージョン, 参照テーブルバージョン] 内の増分変更を返します。
消費せずに読み取りを行う:単独で SELECT 文を実行しても、オフセットバージョンは進捗しません。変更は可視化されますが、消費済みとしてマークされないため、必要に応じて何度でも再読み取り可能です。
消費を行う:ストリームを DML 文(例: INSERT INTO ... SELECT ... FROM <stream_name>)内で使用すると、オフセットバージョンは参照テーブルバージョンに一致するまで進捗します。消費後は、新しい変更が到着するまでストリームは空の結果を返します。
読み取りモードの選択
ストリームを作成する際に read_mode を設定することで、ストリームが返す内容を制御できます。
| モード | 返される内容 | 推奨用途 |
|---|---|---|
append | 変更された各行の最終状態。削除された行は含まれません | 挿入または更新されたデータのみを処理する標準的な ETL パイプライン |
cdc | すべての変更状態(UPDATE 前後の状態、INSERT、DELETE)および 3 つのシステム列 | リアルタイム同期や監査パイプラインなど、完全な変更履歴を必要とする送信先システム |
ストリームの作成
CREATE STREAM [IF NOT EXISTS] <stream_name>
ON TABLE <delta_table_name> <TIMESTAMP AS OF t | VERSION AS OF v>
strmproperties ("read_mode"="append" | "cdc")
[comment <stream_comment>];| パラメーター | 必須 | 説明 |
|---|---|---|
IF NOT EXISTS | いいえ | 省略した場合、同名のストリームが存在するとエラーが返されます。指定した場合は、同名のストリームが存在していてもステートメントは成功し、既存ストリームのメタデータは変更されません。 |
stream_name | はい | 作成するストリームの名前です。 |
ON TABLE <delta_table_name> | はい | ストリームに関連付けるソース Delta テーブルです。ストリームは 1 つのソーステーブルのみをサポートし、作成後にソーステーブルを変更することはできません。 |
TIMESTAMP AS OF t | いいえ | 初期オフセットバージョンをタイムスタンプ t に設定します。クエリ範囲は (t, 最新の増分データタイムスタンプ] から開始されます。 |
VERSION AS OF v | いいえ | 初期オフセットバージョンをデータバージョン v に設定します。クエリ範囲は (v, 最新の増分データバージョン] から開始されます。 |
strmproperties | はい | 文字列形式のキーと値のペアで指定するストリームプロパティです。現在は read_mode のみがサポートされています。有効な値は append および cdc です。 |
stream_comment | いいえ | ストリームのコメントです。最大 1024 バイトまで。これを超えるとエラーが返されます。 |
CDC システム列
read_mode を cdc に設定すると、出力行の末尾に 3 つのシステム列が追加されます。
| 列 | 型 | 説明 |
|---|---|---|
__meta_timestamp | timestamp | 変更が Delta テーブルに書き込まれた時刻です。 |
__meta_op_type | tinyint | 操作タイプ: INSERT(1)または DELETE(0)です。 |
__meta_is_update | tinyint | 行が UPDATE 操作の一部であるかどうか: TRUE(1)または FALSE(0)です。 |
UPDATE は DELETE/INSERT のペアとして表現されるため、__meta_op_type と __meta_is_update を組み合わせることで、正確な変更タイプを特定できます。
| 操作 | __meta_op_type | __meta_is_update |
|---|---|---|
| 新規挿入 | INSERT(1) | FALSE(0) |
| UPDATE 後の値 | INSERT(1) | TRUE(1) |
| UPDATE 前の値 | DELETE(0) | TRUE(1) |
| 削除 | DELETE(0) | FALSE(0) |
例
Delta テーブルを作成し、バージョン 1 から開始する append モードのストリームを作成します。
CREATE TABLE delta_table_src (
pk bigint NOT NULL PRIMARY KEY,
val bigint
) tblproperties ("transactional"="true");
CREATE STREAM delta_table_stream
ON TABLE delta_table_src VERSION AS OF 1
strmproperties('read_mode'='append')
comment 'Stream demo';ストリーム情報の表示
DESC STREAM <stream_name>;例
CREATE TABLE delta_table_src (pk BIGINT NOT NULL PRIMARY KEY,
val BIGINT) TBLPROPERTIES ("transactional"="true");
CREATE STREAM delta_table_stream ON TABLE delta_table_src
VERSION AS OF 1 strmproperties('read_mode'='append')
comment 'Stream demo';
DESC STREAM delta_table_stream;出力例:
Name delta_table_stream
Project sql_optimizer
Create Time 2024-09-06 17:03:32
Last Modified Time 2024-09-06 17:03:32
Offset Version 1
Reference Table Project sql_optimizer
Reference Table Name delta_table_src
Reference Table Id 5e19a67eb97b4477b7fbce0c7bbcebca
Reference Table Version 1
Parameters {
"comment": "stream demo",
"read_mode": "append"}| フィールド | 説明 |
|---|---|
Name | ストリームの名前です。 |
Project | ストリームが存在するプロジェクトです。 |
Create Time | ストリームが作成された時刻です。 |
Last Modified Time | ストリームが最後に変更された時刻です。 |
Offset Version | このストリームによって消費済みのデータバージョンです。 |
Reference Table Project | 関連付けられたソーステーブルが存在するプロジェクトです。 |
Reference Table Name | 関連付けられたソーステーブルの名前です。 |
Reference Table Id | 関連付けられたソーステーブルの固有 ID です。 |
Reference Table Version | 関連付けられたソーステーブルの最新のデータバージョンです。 |
Parameters | ストリームプロパティ(comment および read_mode を含む)です。 |
空のテーブルに対してストリームを初めて作成した場合、Offset VersionとReference Table Versionは等しくなります。Delta テーブルに対して DML 操作が実行されると、Reference Table Versionが進捗します。ストリームは(Offset Version, Reference Table Version]内のすべての変更を返します。DML を用いた読み取りによりこれらの変更が消費されると、Offset VersionがReference Table Versionに追いつき、新しい変更が到着するまでストリームは空の結果を返します。
ストリームの変更
ストリームプロパティの変更
ALTER STREAM <stream_name> SET strmproperties ("key"="value");現在、ストリーム作成後に read_mode を変更することはできません。初期データバージョンの変更
このコマンドを使用してオフセットバージョンをリセットします。たとえば、過去の変更範囲をスキップし、開始ポイントを進めることができます。
ALTER STREAM <stream_name> ON TABLE <delta_table_name>
<TIMESTAMP AS OF t | VERSION AS OF v>;| パラメーター | 説明 |
|---|---|
stream_name | 変更するストリームの名前です。 |
ON TABLE <delta_table_name> | 元のストリームと同じソーステーブルである必要があります。ソーステーブルの変更はサポートされていません。 |
TIMESTAMP AS OF t | オフセットバージョンをタイムスタンプ t にリセットします。クエリ範囲は (t, 最新の増分データバージョン] になります。 |
VERSION AS OF v | オフセットバージョンをデータバージョン v にリセットします。クエリ範囲は (v, 最新の増分データバージョン] になります。 |
例
この例では、ストリームのライフサイクル全体を示します:ストリームの作成、Delta テーブルのバージョンを進めるためのデータ挿入、そしてストリームのオフセットバージョンのリセットです。
-- 1. ソース Delta テーブルを作成します。
CREATE TABLE delta_table_src (pk bigint NOT NULL PRIMARY KEY,
val bigint) tblproperties ("transactional"="true");
-- 2. バージョン 1 から開始するストリームを作成します。
CREATE STREAM delta_table_stream ON TABLE delta_table_src
VERSION AS OF 1 strmproperties('read_mode'='append')
comment 'Stream demo';
-- 3. オフセットバージョンと参照テーブルバージョンがともに 1 であることを確認します。
DESC STREAM delta_table_stream;
-- 出力例:
-- Offset Version 1
-- Reference Table Version 1
-- 4. Delta テーブルのバージョンを新しいものに進めるためにレコードを挿入します。
INSERT INTO delta_table_src VALUES ('1', '1');
-- 5. Delta テーブルのバージョン履歴を表示します。
SHOW HISTORY FOR TABLE delta_table_src;
-- ObjectType ObjectId ObjectName VERSION(LSN) Time Operation
-- TABLE 8605276ce0034b20af761bf4761ba62e delta_table_src 0000000000000001 2024-09-07 10:25:59 CREATE
-- TABLE 8605276ce0034b20af761bf4761ba62e delta_table_src 0000000000000002 2024-09-07 10:28:19 APPEND
-- 6. ストリームのオフセットバージョンをバージョン 2 にリセットし、
-- ステップ 4 で挿入されたデータをスキップします。
ALTER STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 2;
-- 7. 両方のバージョンが現在 2 であることを確認します。
DESC STREAM delta_table_stream;
-- 出力例:
-- Offset Version 2
-- Reference Table Version 2プロジェクト内のすべてのストリームを一覧表示
-- 1. ストリームが存在することを確認します。
SHOW STREAMS;
-- 出力:
-- delta_table_stream
-- 2. ストリームを削除します。
DROP STREAM IF EXISTS delta_table_stream;
-- 3. ストリームが削除されたことを確認します。
SHOW STREAMS;
-- 出力: (empty)例
-- 現在のプロジェクト内のすべてのストリームを一覧表示します。
SHOW STREAMS;
-- 出力例:
-- delta_table_streamストリームの削除
DROP STREAM [IF EXISTS] <stream_name>;例
-- 1. ストリームが存在することを確認します。
SHOW STREAMS;
-- 出力例:
-- delta_table_stream
-- 2. ストリームを削除します。
DROP STREAM IF EXISTS delta_table_stream;
-- 3. ストリームが削除されたことを確認します。
SHOW STREAMS;
-- 出力例:(空)ストリームのクエリ
INSERT INTO <destination_table> SELECT * FROM <stream_name>;DML 文内で使用して変更を消費し、オフセットバージョンを進めるには、以下のように記述します。
INSERT INTO <destination_table> SELECT * FROM <stream_name>;例:CDC モード
この例では、ソーステーブルの INSERT および UPDATE を追跡し、CDC モードを使用して変更を送信先テーブルにコピーします。
Delta テーブルにおける CDC モードは招待プレビュー機能です。設定方法の詳細については、「CDC(招待プレビュー)」をご参照ください。
CDC を有効化したソース Delta テーブルを作成します。
CREATE TABLE delta_table_src ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ( "transactional"="true", 'acid.cdc.mode.enable'='true', 'cdc.insert.into.passthrough.enable'='true' );送信先テーブルを作成します。
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true");CDC モードのストリームを作成します。
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties('read_mode'='cdc') comment 'Stream cdc mode';ソーステーブルに 2 件のレコードを挿入します。
INSERT INTO delta_table_src VALUES (1, 1), (2, 2);ストリームをクエリします。
SELECT文を単独で実行してもオフセットバージョンは進捗しないため、以降の実行でも同じ結果が返されます。SELECT * FROM delta_table_stream; -- 出力例 +------------+------------+------------------+----------------+------------------+ | pk | val | __meta_timestamp | __meta_op_type | __meta_is_update | +------------+------------+------------------+----------------+------------------+ | 2 | 2 | 2024-09-07 11:03:53 | 1 | 0 | | 1 | 1 | 2024-09-07 11:03:53 | 1 | 0 | +------------+------------+------------------+----------------+------------------+両方の行で
__meta_op_type=1(INSERT)および__meta_is_update=0(FALSE)が表示されており、これは新規挿入であることを示しています。送信先テーブルへの挿入により変更を消費します。これによりオフセットバージョンが進捗します。
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream;送信先テーブルにデータが正しく受信されたことを確認します。
SELECT * FROM delta_table_dest; -- 出力例 +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+再度ストリームをクエリします。ステップ 6 で変更が消費されたため、空の結果が返されます。
SELECT * FROM delta_table_stream; -- 出力例 +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ソーステーブルの
pk=1を更新します。UPDATE delta_table_src SET val = 10 WHERE pk = 1;再度ストリームをクエリします。UPDATE は、更新前の状態と更新後の状態という 2 行で表示されます。
SELECT * FROM delta_table_stream; -- 出力例 +------------+------------+------------------+----------------+------------------+ | pk | val | __meta_timestamp | __meta_op_type | __meta_is_update | +------------+------------+------------------+----------------+------------------+ | 1 | 1 | 2024-09-07 11:10:21 | 0 | 1 | | 1 | 10 | 2024-09-07 11:10:21 | 1 | 1 | +------------+------------+------------------+----------------+------------------+1 行目(
__meta_op_type=0、__meta_is_update=1)は更新前の値(DELETE + TRUE = UPDATE_BEFORE)であり、2 行目(__meta_op_type=1、__meta_is_update=1)は更新後の値(INSERT + TRUE = UPDATE_AFTER)です。
例:Append モード
この例では、UPDATE および DELETE 操作に対する Append モードと CDC モードの動作の違いを示します。
ソース Delta テーブルを作成します。
CREATE TABLE delta_table_src ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true");送信先テーブルを作成します。
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true");Append モードのストリームを作成します。
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties ('read_mode'='append') comment 'Stream append mode';ソーステーブルに 2 件のレコードを挿入します。
INSERT INTO delta_table_src VALUES (1, 1), (2, 2);ストリームをクエリします。Append モードではシステム列は返されません。
SELECT * FROM delta_table_stream; -- 出力: +------------+------------+ | pk | val | +------------+------------+ +------------+------------+変更内容を処理する。
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream;送信先テーブルにデータが正しく受信されたことを確認します。
SELECT * FROM delta_table_dest; -- 出力例 +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ストリームをクエリします。ステップ 6 の変更が消費されたため、空の結果が返されます。
SELECT * FROM delta_table_stream; -- 出力例 +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ソーステーブルで
pk=1を更新し、pk=2を削除します。UPDATE delta_table_src SET val = 10 WHERE pk = 1; DELETE FROM delta_table_src WHERE pk = 2;ストリームをクエリします。
SELECT * FROM delta_table_stream; -- 出力例 +------------+------------+ | pk | val | +------------+------------+ | 1 | 10 | +------------+------------+更新された行
(1, 10)のみが返されます。削除された行は含まれません。Append モードは、変更された行の最終状態のみを返し、変更前イメージや削除操作を公開しません。継続的に挿入または更新されるデータを処理する ETL パイプラインには Append モードを、削除や更新前の値を含む完全な変更履歴を必要とする送信先システムには CDC モードを使用してください。
注意事項
各ストリームは、正確に 1 つのソース Delta テーブルのみを追跡します。作成後のソーステーブルの変更はサポートされていません。
read_modeは、ストリーム作成後に変更できません。SELECT文を単独で実行してもオフセットバージョンは進捗しません。変更を消費してポインターを進めるのは、INSERT INTO ... SELECT ... FROM <stream_name>などの DML 操作のみです。異なる送信先システムが同一の変更データを独立して必要とするマルチコンシューマーパイプラインでは、各コンシューマーごとに個別のストリームを作成してください。ストリームはデータを保存せず、バージョンポインターのみを保持するため、同一の Delta テーブルに対する複数のストリームはサポートされています。
ストリームのコメントは 1024 バイトを超えてはなりません。