A função TUMBLE atribui cada linha de um fluxo de dados a uma janela fixa e sem sobreposição. Use-a em uma cláusula GROUP BY para calcular agregações por janela com base no tempo de evento ou no tempo de processamento.
Por exemplo, uma janela tumbling de 5 minutos aplicada a um fluxo de dados infinito gera intervalos discretos e consecutivos: [0:00, 0:05), [0:05, 0:10), [0:10, 0:15), e assim por diante.
Sintaxe
TUMBLE(<time-attr>, <size-interval>)
<size-interval>: INTERVAL 'string' timeUnit
O parâmetro <time-attr> deve ser um atributo de tempo válido: tempo de evento ou tempo de processamento. Para definir atributos de tempo, consulte Visão geral. Para a especificação completa, acesse a documentação do Apache Flink sobre Time Attributes.
Funções identificadoras de janela
Use estas funções na lista SELECT para obter os limites da janela ou propagar o atributo de tempo para janelas em cascata.
|
Função |
Tipo de retorno |
Descrição |
|
|
TIMESTAMP |
Início da janela (inclusivo). Para |
|
|
TIMESTAMP |
Fim da janela (inclusivo). Para |
|
|
TIMESTAMP (atributo rowtime) |
Fim da janela (exclusivo). Para |
|
|
TIMESTAMP (atributo rowtime) |
Fim da janela (exclusivo). Para |
Exemplo 1: Contar cliques por usuário por minuto (tempo de evento)
Este exemplo conta quantas vezes cada usuário clica em uma URL dentro de cada janela tumbling de 1 minuto, usando tempo de evento com uma marca d'água de 2 segundos.
Dados de teste
|
username (VARCHAR) |
click_url (VARCHAR) |
eventtime (VARCHAR) |
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Timo |
|
|
|
Timo |
|
|
Instruções de teste
CREATE TEMPORARY TABLE user_clicks (
username VARCHAR,
click_url VARCHAR,
eventtime VARCHAR,
ts AS TO_TIMESTAMP(eventtime),
WATERMARK FOR ts AS ts - INTERVAL '2' SECOND -- Define a watermark for the rowtime.
) WITH (
'connector' = 'sls',
...
);
CREATE TEMPORARY TABLE tumble_output (
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) WITH (
'connector' = 'datahub',
...
);
INSERT INTO tumble_output
SELECT
TUMBLE_START(ts, INTERVAL '1' MINUTE) AS window_start,
TUMBLE_END(ts, INTERVAL '1' MINUTE) AS window_end,
username,
COUNT(click_url)
FROM user_clicks
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE), username;
Resultados do teste
|
window_start (TIMESTAMP) |
window_end (TIMESTAMP) |
username (VARCHAR) |
clicks (BIGINT) |
|
|
|
Jark |
3 |
|
|
|
Jark |
2 |
|
|
|
Timo |
1 |
Exemplo 2: Contar cliques por usuário por minuto (tempo de processamento)
Este exemplo usa tempo de processamento em vez de tempo de evento. A função PROCTIME() gera automaticamente a coluna de tempo de processamento, eliminando a necessidade de marca d'água.
O Simple Log Service (SLS) aceita apenas colunas do tipo VARCHAR em instruções DDL; portanto, este exemplo grava a saída no DataHub.
Dados de teste
|
username (VARCHAR) |
click_url (VARCHAR) |
|
Jark |
|
|
Jark |
|
|
Jark |
|
|
Jark |
|
|
Jark |
|
|
Timo |
|
Instruções de teste
CREATE TEMPORARY TABLE window_test (
username VARCHAR,
click_url VARCHAR,
ts AS PROCTIME()
) WITH (
'connector' = 'sls',
...
);
CREATE TEMPORARY TABLE tumble_output (
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) WITH (
'connector' = 'datahub',
...
);
INSERT INTO tumble_output
SELECT
TUMBLE_START(ts, INTERVAL '1' MINUTE),
TUMBLE_END(ts, INTERVAL '1' MINUTE),
username,
COUNT(click_url)
FROM window_test
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE), username;
Resultados do teste
|
window_start (TIMESTAMP) |
window_end (TIMESTAMP) |
username (VARCHAR) |
clicks (BIGINT) |
|
|
|
Jark |
5 |
|
|
|
Timo |
1 |
Próximos passos
Visão geral — saiba mais sobre atributos de tempo e janelas em cascata