A janela OVER é uma função de janelamento padrão em bancos de dados tradicionais. Diferentemente da janela GROUP BY, ela calcula um agregado sobre um intervalo de linhas relativo à linha atual sem colapsar o conjunto de resultados. Cada linha mantém sua própria saída, e os elementos do stream podem existir em múltiplas janelas. Cada elemento aciona exatamente um cálculo, e a linha que dispara esse cálculo corresponde à última linha da janela desse elemento.
O Realtime Compute for Apache Flink gerencia todos os dados da janela OVER centralmente em uma única cópia. Logicamente, o sistema cria uma janela separada para cada elemento; após a conclusão do cálculo, descarta os dados desnecessários. Para obter mais informações, consulte Over Aggregation.
Sintaxe
SELECT
agg1(col1) OVER (definition1) AS colName,
...
aggN(colN) OVER (definition1) AS colNameN
FROM Tab1;
agg1(col1): Função de agregação aplicada à colunacol1.OVER (definition1): Cláusula OVER que define a janela.AS colName: Alias da coluna de resultado, referenciável em consultas externas.
A definição OVER deve ser idêntica para todas as agregações, de agg1 a aggN.
Tipos de janela
As janelas OVER dividem-se em dois tipos, conforme a definição do limite de cálculo:
|
Tipo |
Limite de cálculo |
|
Janela ROWS OVER |
Cada linha estabelece seu próprio limite. |
|
Janela RANGE OVER |
Todas as linhas com o mesmo valor de timestamp compartilham um limite. |
Ambos os tipos oferecem suporte a proctime e eventtime.
Janela ROWS OVER
Na janela ROWS OVER, cada linha funciona como um limite de cálculo distinto. Assim, o frame da janela contém uma quantidade fixa de linhas anteriores.
Sintaxe
SELECT
agg1(col1) OVER (
[PARTITION BY value_expression1, ..., value_expressionN]
ORDER BY timeCol
ROWS BETWEEN (UNBOUNDED | rowCount) PRECEDING AND CURRENT ROW
) AS colName,
...
FROM Tab1;
|
Parâmetro |
Descrição |
|
|
Expressão de particionamento. |
|
|
Campo de atributo de tempo para ordenação dos elementos. |
|
|
Quantidade de linhas anteriores a incluir no frame da janela, em relação à linha atual. |
Exemplo
Este exemplo identifica o preço mais alto entre o produto atual e os dois produtos anteriores do mesmo tipo, usando ROWS BETWEEN 2 PRECEDING AND CURRENT ROW.
Tabela de entrada: tmall_item
|
itemid (VARCHAR) |
itemtype (VARCHAR) |
eventtime (VARCHAR) |
price (DOUBLE) |
|
ITEM001 |
Electronic |
|
20 |
|
ITEM002 |
Electronic |
|
50 |
|
ITEM003 |
Electronic |
|
30 |
|
ITEM004 |
Electronic |
|
60 |
|
ITEM005 |
Electronic |
|
40 |
|
ITEM006 |
Electronic |
|
20 |
|
ITEM007 |
Electronic |
|
70 |
|
ITEM008 |
Clothes |
|
20 |
Código de exemplo
CREATE TEMPORARY TABLE tmall_item(
itemid VARCHAR,
itemtype VARCHAR,
eventtime VARCHAR,
onselltime AS TO_TIMESTAMP(eventtime),
price DOUBLE,
WATERMARK FOR onselltime AS onselltime - 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'
);
SELECT
itemid,
itemtype,
onselltime,
price,
MAX(price) OVER (
PARTITION BY itemtype
ORDER BY onselltime
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW
) AS maxprice
FROM tmall_item;
Resultados
|
itemid |
itemtype |
onselltime |
price |
maxprice |
|
ITEM001 |
Electronic |
|
20 |
20 |
|
ITEM002 |
Electronic |
|
50 |
50 |
|
ITEM003 |
Electronic |
|
30 |
50 |
|
ITEM004 |
Electronic |
|
60 |
60 |
|
ITEM005 |
Electronic |
|
40 |
60 |
|
ITEM006 |
Electronic |
|
20 |
60 |
|
ITEM007 |
Electronic |
|
70 |
70 |
|
ITEM008 |
Clothes |
|
20 |
20 |
Durante a fase de aquecimento, quando há menos linhas disponíveis do que o tamanho solicitado pela janela, ocorre uma redução automática da janela. Para ITEM001 (única linha Electronic até aquele momento), a janela contém apenas essa linha; já para ITEM002, ela abrange duas linhas. A janela atinge seu tamanho completo de três linhas a partir de ITEM003.
Janela RANGE OVER
Em uma janela RANGE OVER, todas as linhas com o mesmo valor na coluna de ordenação compartilham um limite de cálculo. O frame da janela é definido por um intervalo de tempo, e não por uma contagem de linhas, o que permite variação na quantidade de linhas em cada frame.
Sintaxe
SELECT
agg1(col1) OVER (
[PARTITION BY value_expression1, ..., value_expressionN]
ORDER BY timeCol
RANGE BETWEEN (UNBOUNDED | timeInterval) PRECEDING AND CURRENT ROW
) AS colName,
...
FROM Tab1;
|
Parâmetro |
Descrição |
|
|
Expressão de particionamento. |
|
|
Campo de atributo de tempo para ordenação dos elementos. |
|
|
Intervalo de tempo das linhas anteriores a incluir no frame da janela, em relação à linha atual. |
Exemplo
Este exemplo busca o preço mais alto entre produtos do mesmo tipo listados nos dois minutos anteriores ao horário de listagem do produto atual, usando RANGE BETWEEN INTERVAL '2' MINUTE PRECEDING AND CURRENT ROW.
A tabela de entrada é a mesma tabela tmall_item do exemplo anterior de janela ROWS OVER.
Código de exemplo
CREATE TEMPORARY TABLE tmall_item(
itemid VARCHAR,
itemtype VARCHAR,
eventtime VARCHAR,
onselltime AS TO_TIMESTAMP(eventtime),
price DOUBLE,
WATERMARK FOR onselltime AS onselltime - 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'
);
SELECT
itemid,
itemtype,
onselltime,
price,
MAX(price) OVER (
PARTITION BY itemtype
ORDER BY onselltime
RANGE BETWEEN INTERVAL '2' MINUTE PRECEDING AND CURRENT ROW
) AS maxprice
FROM tmall_item;
Resultados
|
itemid |
itemtype |
onselltime |
price |
maxprice |
|
ITEM001 |
Electronic |
|
20 |
20 |
|
ITEM002 |
Electronic |
|
50 |
50 |
|
ITEM003 |
Electronic |
|
30 |
50 |
|
ITEM004 |
Electronic |
|
60 |
60 |
|
ITEM005 |
Electronic |
|
40 |
60 |
|
ITEM006 |
Electronic |
|
20 |
40 |
|
ITEM007 |
Electronic |
|
70 |
70 |
|
ITEM008 |
Clothes |
|
20 |
20 |
Observe que o valor de maxprice para ITEM006 é 40, e não 60. Como ITEM004 foi listado às 10:03, ou seja, mais de dois minutos antes do horário de listagem de ITEM006 (10:06), ele fica fora da janela de ITEM006. Apenas ITEM005 (10:05) e ITEM006 (10:06) estão dentro do intervalo de dois minutos.
Essa é a principal diferença comportamental em relação à janela ROWS OVER: uma janela RANGE utiliza a distância temporal, e não a contagem de linhas, para determinar quais dados entram no frame.