Use a função HOP na cláusula GROUP BY para definir uma janela deslizante (hopping window) no Flink SQL. Cada evento pode pertencer simultaneamente a várias janelas sobrepostas, o que torna a função HOP ideal para calcular agregações móveis — por exemplo, contar cliques por usuário no último minuto, com recálculo a cada 30 segundos.
Sintaxe
HOP(<time-attr>, <slide-interval>, <size-interval>)
Parâmetros
|
Parâmetro |
Descrição |
Exemplo |
|
|
Coluna de atributo de tempo no fluxo. Define se o sistema usará tempo de evento ou tempo de processamento. Para mais informações, consulte Atributos de tempo. |
— |
|
|
Frequência de avanço da janela. Determina a diferença de tempo entre o início de janelas consecutivas. Formato: |
|
|
|
Duração total abrangida por cada janela. Formato: |
|
Comportamento de sobreposição de janelas
A relação entre <slide-interval> e <size-interval> determina como as janelas se sobrepõem:
|
Condição |
Comportamento |
|
|
As janelas se sobrepõem. O sistema atribui cada evento a múltiplas janelas. Este é o padrão clássico de janela deslizante. |
|
|
As janelas são contíguas, sem intervalos — equivalente a janelas saltitantes (tumbling windows). |
|
|
Não há sobreposição; as janelas ficam separadas por lacunas temporais. |
Funções identificadoras de janela
Use estas funções na cláusula SELECT para obter o horário inicial, final ou o atributo de tempo de uma janela.
|
Função |
Tipo de retorno |
Descrição |
|
|
TIMESTAMP |
Horário inicial da janela (inclusivo). Por exemplo, retorna |
|
|
TIMESTAMP |
Horário final da janela (inclusivo). Por exemplo, retorna |
|
|
TIMESTAMP (rowtime-attr) |
Horário final da janela (exclusivo). Por exemplo, retorna |
|
|
TIMESTAMP (rowtime-attr) |
Horário final da janela (exclusivo). Por exemplo, retorna |
Exemplo
Este exemplo conta cliques por usuário em uma janela deslizante de 1 minuto que avança a cada 30 segundos. A consulta usa HOP_START e HOP_END para recuperar os limites da janela.
Dados de teste (user_clicks)
|
username (VARCHAR) |
click_url (VARCHAR) |
eventtime (VARCHAR) |
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Timo |
|
|
SQL
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' = 'kafka',
'topic' = '<yourTopic>',
'properties.bootstrap.servers' = '<brokers>',
'scan.startup.mode' = 'earliest-offset',
'format' = 'csv'
);
CREATE TEMPORARY TABLE hop_output (
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO hop_output
SELECT
HOP_START(ts, INTERVAL '30' SECOND, INTERVAL '1' MINUTE),
HOP_END(ts, INTERVAL '30' SECOND, INTERVAL '1' MINUTE),
username,
COUNT(click_url)
FROM user_clicks
GROUP BY HOP(ts, INTERVAL '30' SECOND, INTERVAL '1' MINUTE), username;
Resultados
|
window_start (TIMESTAMP) |
window_end (TIMESTAMP) |
username (VARCHAR) |
clicks (BIGINT) |
|
|
|
Jark |
2 |
|
|
|
Jark |
3 |
|
|
|
Jark |
2 |
|
|
|
Jark |
2 |
|
|
|
Jark |
1 |
|
|
|
Timo |
1 |
① Horário inicial da primeira janela
Se a janela deslizante não conseguir determinar o momento exato em que o primeiro evento entrou no fluxo, ela desloca o início da primeira janela para trás conforme a fórmula: duração da janela - passo de deslizamento.
Veja o exemplo:
|
Duração da janela (segundos) |
Passo de deslizamento (segundos) |
Tempo do evento |
Início da primeira janela |
Fim da primeira janela |
|
120 |
30 |
|
|
|
|
60 |
10 |
|
|
|
② Janela ainda não disparada
A linha com window_end = 2024-10-10 10:02:30.0 não aparece nos resultados porque a janela ainda não foi acionada. O disparo ocorre quando:
event time >= window_end + watermark offset
Por exemplo: 10:02:30.0 + 2 seconds = 10:02:32.0. Qualquer evento de qualquer usuário ocorrido em ou após 10:02:32.0 aciona essa janela.