Un stream est un objet MaxCompute qui gère automatiquement les versions de données pour les requêtes incrémentielles sur les tables Delta. Il suit les modifications du langage de manipulation de données (DML) — insertions, mises à jour et suppressions — ainsi que les métadonnées associées à chaque changement, et les rend disponibles pour une consommation incrémentielle. Chaque stream conserve un pointeur de version afin que les consommateurs sachent toujours quels changements ont été traités et lesquels sont nouveaux.
Cette rubrique présente les commandes SQL permettant de créer, inspecter, modifier, lister, supprimer et interroger des streams.
Fonctionnement
Un stream est toujours associé à une seule table Delta. En interne, il maintient deux marqueurs de version :
Offset Version : la version de données jusqu'à laquelle les modifications ont été consommées. Cette valeur n'avance que lors de la lecture du stream dans le cadre d'une opération DML.
Reference Table Version : la dernière version de données de la table Delta associée. Cette valeur se met à jour automatiquement à mesure que la table évolue.
Chaque interrogation d'un stream renvoie les modifications incrémentielles situées dans l'intervalle semi-ouvert (Offset Version, Reference Table Version].
Lecture sans consommation : L'exécution d'une instruction SELECT seule ne fait pas avancer l'Offset Version. Les modifications restent visibles mais ne sont pas marquées comme consommées ; vous pouvez les relire autant de fois que nécessaire.
Consommation : Lorsque vous utilisez un stream dans une instruction DML (par exemple, INSERT INTO ... SELECT ... FROM <stream_name>), l'Offset Version avance pour correspondre à la Reference Table Version. Après la consommation, le stream renvoie un résultat vide jusqu'à l'arrivée de nouvelles modifications.
Choisir un mode de lecture
Définissez read_mode lors de la création d'un stream pour contrôler ce que celui-ci renvoie.
| Mode | Ce qu'il renvoie | Idéal pour |
|---|---|---|
append |
État final de chaque ligne modifiée ; les lignes supprimées sont exclues | Pipelines ETL standards qui traitent uniquement les données insérées ou mises à jour |
cdc |
Tous les états de modification (avant et après mise à jour, insertion, suppression), plus trois colonnes système | Systèmes en aval nécessitant l'historique complet des modifications, tels que la synchronisation en temps réel ou les pipelines d'audit |
Créer un stream
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>];
| Paramètre | Obligatoire | Description |
|---|---|---|
IF NOT EXISTS |
Non | Si omis et qu'un stream du même nom existe, une erreur est renvoyée. Si spécifié, l'instruction réussit même si un stream portant le même nom existe ; les métadonnées du stream existant restent inchangées. |
stream_name |
Oui | Nom du stream à créer. |
ON TABLE <delta_table_name> |
Oui | Table Delta source à associer au stream. Un stream ne prend en charge qu'une seule table source, et la table source ne peut pas être modifiée après la création. |
TIMESTAMP AS OF t |
Non | Définit l'Offset Version initiale sur l'horodatage t. La plage de requête commence à partir de (t, horodatage des dernières données incrémentielles]. |
VERSION AS OF v |
Non | Définit l'Offset Version initiale sur la version de données v. La plage de requête commence à partir de (v, dernière version de données incrémentielles]. |
strmproperties |
Oui | Propriétés du stream sous forme de paires clé-valeur de type chaîne. Actuellement, seul read_mode est pris en charge. Valeurs valides : append et cdc. |
stream_comment |
Non | Commentaire pour le stream. Maximum 1024 octets ; une erreur est renvoyée en cas de dépassement. |
Colonnes système CDC
Lorsque read_mode est défini sur cdc, trois colonnes système sont ajoutées à chaque ligne de sortie :
| Colonne | Type | Description |
|---|---|---|
__meta_timestamp |
timestamp | Moment où la modification a été écrite dans la table Delta. |
__meta_op_type |
tinyint | Type d'opération : INSERT (1) ou DELETE (0). |
__meta_is_update |
tinyint | Indique si la ligne fait partie d'une mise à jour : TRUE (1) ou FALSE (0). |
Étant donné que les mises à jour sont représentées par une paire DELETE/INSERT, la combinaison de __meta_op_type et de __meta_is_update permet d'identifier le type exact de modification :
| Opération | **__meta_op_type** |
**__meta_is_update** |
|---|---|---|
| Nouvelle insertion | INSERT (1) | FALSE (0) |
| Valeur après une mise à jour | INSERT (1) | TRUE (1) |
| Valeur avant une mise à jour | DELETE (0) | TRUE (1) |
| Suppression | DELETE (0) | FALSE (0) |
Exemple
Créez une table Delta, puis créez un stream en mode append à partir de la version 1.
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';
Afficher les informations d'un stream
DESC STREAM <stream_name>;
Exemple
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;
Sortie :
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"}
| Champ | Description |
|---|---|
Name |
Nom du stream. |
Project |
Projet dans lequel réside le stream. |
Create Time |
Date et heure de création du stream. |
Last Modified Time |
Date et heure de la dernière modification du stream. |
Offset Version |
Version de données jusqu'à laquelle les modifications ont été consommées par ce stream. |
Reference Table Project |
Projet dans lequel réside la table source associée. |
Reference Table Name |
Nom de la table source associée. |
Reference Table Id |
ID unique de la table source associée. |
Reference Table Version |
Dernière version de données de la table source associée. |
Parameters |
Propriétés du stream, y compris comment et read_mode. |
Lors de la création initiale du stream sur une table vide,Offset VersionetReference Table Versionsont identiques. À mesure que des opérations DML s'exécutent sur la table Delta,Reference Table Versionavance. Le stream renvoie toutes les modifications situées dans(Offset Version, Reference Table Version]. Après qu'une lecture basée sur DML a consommé ces modifications,Offset VersionrattrapeReference Table Version, et le stream renvoie un résultat vide jusqu'à l'arrivée de nouvelles modifications.
Modifier un stream
Modifier les propriétés d'un stream
ALTER STREAM <stream_name> SET strmproperties ("key"="value");
Actuellement, read_mode ne peut pas être modifié après la création du stream.
Modifier la version de données initiale
Utilisez cette commande pour réinitialiser l'Offset Version, par exemple pour ignorer une plage de modifications historiques et faire avancer le point de départ.
ALTER STREAM <stream_name> ON TABLE <delta_table_name>
<TIMESTAMP AS OF t | VERSION AS OF v>;
| Paramètre | Description |
|---|---|
stream_name |
Nom du stream à modifier. |
ON TABLE <delta_table_name> |
Doit correspondre à la même table source que l'original. La modification de la table source n'est pas prise en charge. |
TIMESTAMP AS OF t |
Réinitialise l'Offset Version sur l'horodatage t. La plage de requête devient (t, dernière version de données incrémentielles]. |
VERSION AS OF v |
Réinitialise l'Offset Version sur la version de données v. La plage de requête devient (v, dernière version de données incrémentielles]. |
Exemple
Cet exemple illustre le cycle de vie complet : création d'un stream, insertion de données pour faire avancer la version de la table Delta, puis réinitialisation de l'Offset Version du stream.
-- 1. Create a source Delta Table.
CREATE TABLE delta_table_src (pk bigint NOT NULL PRIMARY KEY,
val bigint) tblproperties ("transactional"="true");
-- 2. Create a stream starting from version 1.
CREATE STREAM delta_table_stream ON TABLE delta_table_src
VERSION AS OF 1 strmproperties('read_mode'='append')
comment 'Stream demo';
-- 3. Confirm that Offset Version and Reference Table Version are both 1.
DESC STREAM delta_table_stream;
-- Output:
-- Offset Version 1
-- Reference Table Version 1
-- 4. Insert a record to advance the Delta Table to a new version.
INSERT INTO delta_table_src VALUES ('1', '1');
-- 5. View Delta Table version history.
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. Reset the stream's Offset Version to version 2,
-- skipping the data inserted in step 4.
ALTER STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 2;
-- 7. Confirm that both versions are now 2.
DESC STREAM delta_table_stream;
-- Output:
-- Offset Version 2
-- Reference Table Version 2
Lister tous les streams d'un projet
SHOW STREAMS;
Exemple
-- List all streams in the current project.
SHOW STREAMS;
-- Output:
-- delta_table_stream
Supprimer un stream
DROP STREAM [IF EXISTS] <stream_name>;
Exemple
-- 1. Confirm the stream exists.
SHOW STREAMS;
-- Output:
-- delta_table_stream
-- 2. Delete the stream.
DROP STREAM IF EXISTS delta_table_stream;
-- 3. Confirm the stream is gone.
SHOW STREAMS;
-- Output: (empty)
Interroger un stream
SELECT * FROM <stream_name>;
Utilisez cette instruction au sein d'une instruction DML pour consommer les modifications et faire avancer l'Offset Version :
INSERT INTO <destination_table> SELECT * FROM <stream_name>;
Exemple : mode CDC
Cet exemple suit les insertions et les mises à jour sur une table source et copie les modifications vers une table de destination en utilisant le mode CDC.
Le mode CDC sur les tables Delta nécessite une préversion sur invitation. Pour plus de détails sur la configuration, consultez CDC (préversion sur invitation).
-
Créez une table Delta source avec le CDC activé.
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' ); -
Créez une table de destination.
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Créez un stream en mode CDC.
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties('read_mode'='cdc') comment 'Stream cdc mode'; -
Insérez deux enregistrements dans la table source.
INSERT INTO delta_table_src VALUES (1, 1), (2, 2); -
Interrogez le stream. L'exécution de
SELECTseule ne fait pas avancer l'Offset Version ; le même résultat est renvoyé à chaque exécution suivante.SELECT * FROM delta_table_stream; -- Output +------------+------------+------------------+----------------+------------------+ | 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 | +------------+------------+------------------+----------------+------------------+Les deux lignes affichent
__meta_op_type=1(INSERT) et__meta_is_update=0(FALSE), indiquant de nouvelles insertions. -
Consommez les modifications en les insérant dans la table de destination. Cela fait avancer l'Offset Version.
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream; -
Confirmez que la table de destination a reçu les données.
SELECT * FROM delta_table_dest; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Interrogez à nouveau le stream. Il renvoie un résultat vide car les modifications ont été consommées à l'étape 6.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ -
Mettez à jour
pk=1dans la table source.UPDATE delta_table_src SET val = 10 WHERE pk = 1; -
Interrogez à nouveau le stream. La mise à jour apparaît sous forme de deux lignes : l'état avant la mise à jour et l'état après la mise à jour.
SELECT * FROM delta_table_stream; -- Output +------------+------------+------------------+----------------+------------------+ | 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 | +------------+------------+------------------+----------------+------------------+La première ligne (
__meta_op_type=0,__meta_is_update=1) correspond à la valeur avant la mise à jour (DELETE + TRUE = UPDATE_BEFORE). La deuxième ligne (__meta_op_type=1,__meta_is_update=1) correspond à la valeur après la mise à jour (INSERT + TRUE = UPDATE_AFTER).
Exemple : mode Append
Cet exemple montre la différence de comportement entre le mode append et le mode CDC pour les opérations de mise à jour et de suppression.
-
Créez une table Delta source.
CREATE TABLE delta_table_src ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Créez une table de destination.
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Créez un stream en mode append.
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties ('read_mode'='append') comment 'Stream append mode'; -
Insérez deux enregistrements dans la table source.
INSERT INTO delta_table_src VALUES (1, 1), (2, 2); -
Interrogez le stream. Le mode append ne renvoie aucune colonne système.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Consommez les modifications.
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream; -
Confirmez que la table de destination a reçu les données.
SELECT * FROM delta_table_dest; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Interrogez le stream. Il renvoie un résultat vide — les modifications de l'étape 6 ont été consommées.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ -
Mettez à jour
pk=1et supprimezpk=2dans la table source.UPDATE delta_table_src SET val = 10 WHERE pk = 1; DELETE FROM delta_table_src WHERE pk = 2; -
Interrogez le stream.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 10 | +------------+------------+Seule la ligne mise à jour
(1, 10)est renvoyée. La ligne supprimée n'est pas incluse. Le mode append renvoie uniquement l'état final des lignes modifiées ; il n'expose ni les images avant modification ni les suppressions. Utilisez le mode append pour les pipelines ETL qui traitent des données continuellement insérées ou mises à jour ; utilisez le mode CDC lorsque votre système en aval a besoin de l'historique complet des modifications, y compris les suppressions et les valeurs avant mise à jour.
Notes d'utilisation
Chaque stream suit exactement une table Delta source. La modification de la table source après la création n'est pas prise en charge.
read_modene peut pas être modifié après la création du stream.Une instruction
SELECTseule ne fait pas avancer l'Offset Version ; seules les opérations DML (telles queINSERT INTO ... SELECT ... FROM <stream_name>) consomment les modifications et font avancer le pointeur.Pour les pipelines multi-consommateurs où différents systèmes en aval ont besoin des mêmes données de modification de manière indépendante, créez un stream distinct pour chaque consommateur. Les streams ne stockent pas de données — ils ne conservent qu'un pointeur de version — donc plusieurs streams sur la même table Delta sont pris en charge.
Les commentaires de stream ne doivent pas dépasser 1024 octets.