Um stream é um objeto do MaxCompute que gerencia automaticamente versões de dados para consultas incrementais em Delta Tables. Ele rastreia alterações da linguagem de manipulação de dados (DML) — inserções, atualizações e exclusões — e os metadados correspondentes, disponibilizando-os para consumo incremental. Cada stream mantém um ponteiro de versão para que os consumidores saibam sempre quais alterações já foram processadas e quais são novas.
Este tópico aborda os comandos SQL para criar, inspecionar, modificar, listar, excluir e consultar streams.
Como funciona
Um stream está sempre associado a exatamente uma Delta Table. Internamente, ele mantém dois marcadores de versão:
Offset Version: versão dos dados até a qual as alterações foram consumidas. Esse valor avança apenas quando você lê o stream dentro de uma operação DML.
Reference Table Version: versão mais recente dos dados da Delta Table associada. Este valor se atualiza automaticamente conforme a tabela sofre alterações.
Sempre que você consulta um stream, o MaxCompute retorna as alterações incrementais no intervalo semiaberto (Offset Version, Reference Table Version].
Leitura sem consumo: Executar apenas uma instrução SELECT não avança a Offset Version. As alterações ficam visíveis, mas não são marcadas como consumidas — você pode relê-las quantas vezes forem necessárias.
Consumo: Ao usar um stream dentro de uma instrução DML (por exemplo, INSERT INTO ... SELECT ... FROM <stream_name>), a Offset Version avança para igualar a Reference Table Version. Após o consumo, o stream retorna vazio até que novas alterações cheguem.
Escolha um modo de leitura
Defina read_mode ao criar um stream para controlar o que ele retorna.
|
Modo |
O que retorna |
Mais indicado para |
|
|
Estado final de cada linha alterada; linhas excluídas são omitidas |
Pipelines ETL padrão que processam apenas dados inseridos ou atualizados |
|
|
Todos os estados de alteração (antes e depois da atualização, inserção, exclusão), além de três colunas de sistema |
Sistemas downstream que precisam do histórico completo de alterações, como sincronização em tempo real ou pipelines de auditoria |
Crie um 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>];
|
Parâmetro |
Obrigatório |
Descrição |
|
|
Não |
Se omitido e já existir um stream com o mesmo nome, o sistema retornará um erro. Se especificado, a instrução terá êxito mesmo que exista um stream com o mesmo nome; os metadados do stream existente permanecem inalterados. |
|
|
Sim |
Nome do stream a ser criado. |
|
|
Sim |
Delta Table de source a ser associada ao stream. Um stream suporta apenas uma tabela de source, e essa tabela não pode ser alterada após a criação. |
|
|
Não |
Defina a Offset Version inicial como o timestamp |
|
|
Não |
Defina a Offset Version inicial como a versão de dados |
|
|
Sim |
Propriedades do stream como pares chave-valor em string. Atualmente, apenas |
|
|
Não |
Comentário para o stream. Máximo de 1024 bytes; o sistema retornará um erro se esse limite for excedido. |
Colunas de sistema CDC
Quando read_mode está definido como cdc, três colunas de sistema são adicionadas a cada linha de saída:
|
Coluna |
Tipo |
Descrição |
|
|
timestamp |
Momento em que a alteração foi gravada na Delta Table. |
|
|
tinyint |
Tipo de operação: |
|
|
tinyint |
Indica se a linha faz parte de um UPDATE: |
Como as atualizações são representadas como um par DELETE/INSERT, a combinação de __meta_op_type e __meta_is_update identifica o tipo exato de alteração:
|
Operação |
** |
** |
|
Nova inserção |
INSERT (1) |
FALSE (0) |
|
Valor após uma atualização |
INSERT (1) |
TRUE (1) |
|
Valor antes de uma atualização |
DELETE (0) |
TRUE (1) |
|
Exclusão |
DELETE (0) |
FALSE (0) |
Exemplo
Crie uma Delta Table e, em seguida, crie um stream no modo append começando pela versão 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';
Visualize informações do stream
DESC STREAM <stream_name>;
Exemplo
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;
Saída:
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"}
|
Campo |
Descrição |
|
|
Nome do stream. |
|
|
Projeto onde o stream reside. |
|
|
Horário de criação do stream. |
|
|
Horário da última modificação do stream. |
|
|
Versão dos dados até a qual este stream consumiu alterações. |
|
|
Projeto onde reside a tabela de source associada. |
|
|
Nome da tabela de source associada. |
|
|
ID exclusivo da tabela de source associada. |
|
|
Versão mais recente dos dados da tabela de source associada. |
|
|
Propriedades do stream, incluindo |
Quando o stream é criado pela primeira vez em uma tabela vazia,Offset VersioneReference Table Versionsão iguais. À medida que operações DML são executadas na Delta Table, aReference Table Versionavança. O stream retorna todas as alterações no intervalo(Offset Version, Reference Table Version]. Após uma leitura baseada em DML consumir essas alterações, aOffset Versionalcança aReference Table Version, e o stream retorna vazio até que novas alterações cheguem.
Modifique um stream
Modifique propriedades do stream
ALTER STREAM <stream_name> SET strmproperties ("key"="value");
Atualmente, não é possível modificar o read_mode após a criação do stream.
Modifique a versão inicial dos dados
Use este comando para redefinir a Offset Version — por exemplo, para pular um intervalo de alterações históricas e avançar o ponto de partida.
ALTER STREAM <stream_name> ON TABLE <delta_table_name>
<TIMESTAMP AS OF t | VERSION AS OF v>;
|
Parâmetro |
Descrição |
|
|
Nome do stream a ser modificado. |
|
|
Deve ser a mesma tabela de source da original. Não há suporte para alteração da tabela de source. |
|
|
Redefine a Offset Version para o timestamp |
|
|
Redefine a Offset Version para a versão de dados |
Exemplo
Este exemplo mostra o ciclo de vida completo: crie um stream, insira dados para avançar a versão da Delta Table e, em seguida, redefina a Offset Version do 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
Liste todos os streams em um projeto
SHOW STREAMS;
Exemplo
-- List all streams in the current project.
SHOW STREAMS;
-- Output:
-- delta_table_stream
Exclua um stream
DROP STREAM [IF EXISTS] <stream_name>;
Exemplo
-- 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)
Consulte um stream
SELECT * FROM <stream_name>;
Use isto dentro de uma instrução DML para consumir alterações e avançar a Offset Version:
INSERT INTO <destination_table> SELECT * FROM <stream_name>;
Exemplo: Modo CDC
Este exemplo rastreia inserções e atualizações em uma tabela de source e copia as alterações para uma tabela de destino usando o modo CDC.
O modo CDC em Delta Tables requer preview por convite. Para detalhes de configuração, consulte CDC (preview por convite).
-
Crie uma Delta Table de source com CDC ativado.
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' ); -
Crie uma tabela de destino.
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Crie um stream no modo CDC.
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties('read_mode'='cdc') comment 'Stream cdc mode'; -
Insira dois registros na tabela de source.
INSERT INTO delta_table_src VALUES (1, 1), (2, 2); -
Consulte o stream. Executar apenas
SELECTnão avança a Offset Version — o mesmo resultado é retornado em cada execução subsequente.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 | +------------+------------+------------------+----------------+------------------+Ambas as linhas mostram
__meta_op_type=1(INSERT) e__meta_is_update=0(FALSE), indicando novas inserções. -
Consuma as alterações inserindo-as na tabela de destino. Isso avança a Offset Version.
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream; -
Confirme se a tabela de destino recebeu os dados.
SELECT * FROM delta_table_dest; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Consulte o stream novamente. Ele retorna vazio porque as alterações foram consumidas na etapa 6.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ -
Atualize
pk=1na tabela de source.UPDATE delta_table_src SET val = 10 WHERE pk = 1; -
Consulte o stream novamente. O UPDATE aparece como duas linhas: o estado anterior à atualização e o estado posterior.
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 | +------------+------------+------------------+----------------+------------------+A primeira linha (
__meta_op_type=0,__meta_is_update=1) é o valor antes da atualização (DELETE + TRUE = UPDATE_BEFORE). A segunda linha (__meta_op_type=1,__meta_is_update=1) é o valor após a atualização (INSERT + TRUE = UPDATE_AFTER).
Exemplo: Modo append
Este exemplo mostra a diferença de comportamento entre o modo append e o modo CDC para operações UPDATE e DELETE.
-
Crie uma Delta Table de source.
CREATE TABLE delta_table_src ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Crie uma tabela de destino.
CREATE TABLE delta_table_dest ( pk bigint NOT NULL PRIMARY KEY, val bigint ) tblproperties ("transactional"="true"); -
Crie um stream no modo append.
CREATE STREAM delta_table_stream ON TABLE delta_table_src VERSION AS OF 1 strmproperties ('read_mode'='append') comment 'Stream append mode'; -
Insira dois registros na tabela de source.
INSERT INTO delta_table_src VALUES (1, 1), (2, 2); -
Consulte o stream. O modo append não retorna colunas de sistema.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Consuma as alterações.
INSERT INTO delta_table_dest SELECT pk, val FROM delta_table_stream; -
Confirme se a tabela de destino recebeu os dados.
SELECT * FROM delta_table_dest; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 1 | | 2 | 2 | +------------+------------+ -
Consulte o stream. Ele retorna vazio — as alterações da etapa 6 foram consumidas.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ +------------+------------+ -
Atualize
pk=1e excluapk=2na tabela de source.UPDATE delta_table_src SET val = 10 WHERE pk = 1; DELETE FROM delta_table_src WHERE pk = 2; -
Consulte o stream.
SELECT * FROM delta_table_stream; -- Output +------------+------------+ | pk | val | +------------+------------+ | 1 | 10 | +------------+------------+Apenas a linha atualizada
(1, 10)é retornada. A linha excluída não está incluída. O modo append retorna apenas o estado final das linhas modificadas — ele não expõe imagens anteriores ou exclusões. Use o modo append para pipelines ETL que processam dados inseridos ou atualizados continuamente; use o modo CDC quando seu sistema downstream precisar do histórico completo de alterações, incluindo exclusões e valores anteriores à atualização.
Notas de uso
Cada stream rastreia exatamente uma Delta Table de source. Não há suporte para alteração da tabela de source após a criação.
Não é possível alterar o
read_modeapós a criação do stream.Uma instrução
SELECTsozinha não avança a Offset Version; apenas operações DML (comoINSERT INTO ... SELECT ... FROM <stream_name>) consomem alterações e avançam o ponteiro.Para pipelines com múltiplos consumidores em que diferentes sistemas downstream precisam dos mesmos dados de alteração de forma independente, crie um stream separado para cada consumidor. Streams não armazenam dados — eles armazenam apenas um ponteiro de versão — portanto, há suporte para múltiplos streams na mesma Delta Table.
Os comentários do stream não devem exceder 1024 bytes.