Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:OVER

Última atualização: Jun 27, 2026

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 à coluna col1.

  • 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

value_expression

Expressão de particionamento.

timeCol

Campo de atributo de tempo para ordenação dos elementos.

rowCount

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

2024-11-11 10:01:00

20

ITEM002

Electronic

2024-11-11 10:02:00

50

ITEM003

Electronic

2024-11-11 10:03:00

30

ITEM004

Electronic

2024-11-11 10:03:00

60

ITEM005

Electronic

2024-11-11 10:05:00

40

ITEM006

Electronic

2024-11-11 10:06:00

20

ITEM007

Electronic

2024-11-11 10:07:00

70

ITEM008

Clothes

2024-11-11 10:08:00

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

2024-11-11 10:01:00

20

20

ITEM002

Electronic

2024-11-11 10:02:00

50

50

ITEM003

Electronic

2024-11-11 10:03:00

30

50

ITEM004

Electronic

2024-11-11 10:03:00

60

60

ITEM005

Electronic

2024-11-11 10:05:00

40

60

ITEM006

Electronic

2024-11-11 10:06:00

20

60

ITEM007

Electronic

2024-11-11 10:07:00

70

70

ITEM008

Clothes

2024-11-11 10:08:00

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

value_expression

Expressão de particionamento.

timeCol

Campo de atributo de tempo para ordenação dos elementos.

timeInterval

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

2024-11-11 10:01:00

20

20

ITEM002

Electronic

2024-11-11 10:02:00

50

50

ITEM003

Electronic

2024-11-11 10:03:00

30

50

ITEM004

Electronic

2024-11-11 10:03:00

60

60

ITEM005

Electronic

2024-11-11 10:05:00

40

60

ITEM006

Electronic

2024-11-11 10:06:00

20

40

ITEM007

Electronic

2024-11-11 10:07:00

70

70

ITEM008

Clothes

2024-11-11 10:08:00

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.