A função SESSION agrupa elementos de stream por atividade de sessão e cria janelas de tamanho variável sem sobreposição. A janela de sessão fecha quando nenhum elemento chega durante o período de inatividade definido pelo intervalo da sessão. Elementos recebidos antes do fim desse intervalo são mesclados na mesma sessão; elementos recebidos após o intervalo abrem uma nova janela de sessão.
Por exemplo, com um intervalo de sessão de 10 minutos: se o tempo entre dois eventos do mesmo usuário for inferior a 10 minutos, ambos pertencem à mesma janela de sessão. Se nenhum evento ocorrer nos 10 minutos seguintes ao último evento, a janela fecha e é enviada para downstream. Eventos subsequentes iniciam uma nova janela de sessão.
Sintaxe
Use SESSION em uma cláusula GROUP BY para definir uma janela de sessão.
SESSION(<time-attr>, <gap-interval>)
Parâmetros
|
Parâmetro |
Descrição |
Exemplo |
|
|
Campo de atributo de tempo válido no stream. Define se deve ser usado o tempo de processamento ou o tempo do evento. Para mais detalhes, consulte Time attributes. |
- |
|
|
Intervalo da sessão: período de inatividade após o qual a janela de sessão fecha. Formato: |
|
Funções identificadoras de janela
As funções identificadoras de janela retornam o horário de início, o horário de término ou o atributo de tempo de uma janela de sessão. Use-as em uma cláusula SELECT para obter os limites da janela.
|
Função |
Tipo de retorno |
Descrição |
|
|
TIMESTAMP |
Retorna o horário de início da janela (inclusivo). Por exemplo, para a janela |
|
|
TIMESTAMP |
Retorna o horário de término da janela (inclusivo). Por exemplo, para a janela |
|
|
TIMESTAMP (atributo rowtime) |
Retorna o horário de término da janela (exclusivo). Por exemplo, para a janela |
|
|
TIMESTAMP (atributo de tempo de processamento) |
Retorna o horário de término da janela (exclusivo). Por exemplo, para a janela |
Exemplo
Este exemplo conta cliques por usuário em cada sessão ativa, com um intervalo de sessão de 30 segundos e tempo de evento.
Dados de teste (tabela user_clicks)
|
username (VARCHAR) |
click_url (VARCHAR) |
eventtime (VARCHAR) |
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Jark |
|
|
|
Timo |
|
|
Instruções 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 session_output(
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) WITH (
'connector'='print',
'logger'='true'
);
INSERT INTO session_output
SELECT
SESSION_START(ts, INTERVAL '30' SECOND),
SESSION_END(ts, INTERVAL '30' SECOND),
username,
COUNT(click_url)
FROM user_clicks
GROUP BY SESSION(ts, INTERVAL '30' SECOND), username;
Resultados
|
window_start (TIMESTAMP) |
window_end (TIMESTAMP) |
username (VARCHAR) |
clicks (BIGINT) |
|
|
|
Jark |
2 |
|
|
|
Jark |
2 |
|
|
|
Jark |
1 |
|
|
|
Timo |
1 |