Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Agregação de janela

Última atualização: Jun 27, 2026

O Realtime Compute for Apache Flink oferece suporte a dois tipos de agregação de janela: agregação de janela de grupo e agregação por função com valor de tabela (TVF) de janela. Este tópico aborda a sintaxe de cada uma, os cenários em que a agregação por TVF de janela reverte para o modo não-TVF e o suporte a streams de atualização nos diferentes tipos de janela.

Escolha entre as duas sintaxes

Agregação de janela de grupo

Agregação por TVF de janela

Operador

GroupWindowAggregation

WindowAggregate

Funções de janela

TUMBLE, HOP, SESSION

TUMBLE, HOP, CUMULATE, SESSION

Status

Obsoleto

Recomendado

Otimizações de desempenho

Não

Sim

Suporte a GROUPING SETS

Não

Sim

Window Top-N após agregação

Não

Sim

Suporte a stream de atualização

Sim (VVR)

Sim (VVR, todos os tipos de janela)

Use a agregação por TVF de janela. Ela oferece suporte a todos os tipos de janela da agregação de janela de grupo, além de CUMULATE, inclui otimizações de desempenho e GROUPING SETS, e permite aplicar Window Top-N aos resultados da agregação.

Agregação de janela de grupo (obsoleta)

A agregação de janela de grupo define janelas na cláusula GROUP BY. Ela corresponde ao operador GroupWindowAggregation e oferece suporte às funções de janela TUMBLE, HOP e SESSION.

Para obter detalhes sobre sintaxe, exemplos e recursos, consulte Agregação de janela de grupo.

Mudança de comportamento no VVR 11.x para agregação de janela de grupo

A partir do VVR 11.x (Flink 1.20), o sistema não reescreve mais automaticamente a agregação de janela de grupo (sintaxe obsoleta) em um plano de execução de agregação por TVF de janela.

No VVR 8.x, o sistema reescrevia automaticamente a sintaxe obsoleta para o novo plano de execução, habilitando a otimização de agregação em duas fases Local-Global. A partir do VVR 11.x, essa reescrita automática deixou de ser o comportamento padrão: a sintaxe obsoleta mantém seu plano de execução física original.

Impacto: jobs que utilizam a sintaxe obsoleta não se beneficiam mais da otimização de agregação em duas fases automaticamente. O desempenho pode degradar em cenários com grandes volumes de dados ou skew de dados.

Importante

Migre a agregação de janela de grupo para a agregação por TVF de janela (nova sintaxe). Para conhecer as condições que ativam a otimização de agregação em duas fases, consulte Otimização Local-Global para agregação de janela.

Compatibilidade com o comportamento anterior

Caso a migração imediata não seja viável, ative o seguinte parâmetro para restaurar o comportamento do VVR 8.x:

Parâmetro

Descrição

Padrão

table.optimizer.window-rewrite-enabled

Habilita a reescrita automática da sintaxe obsoleta para a nova sintaxe, ativando a otimização de agregação em duas fases.

false (desde o VVR 11.x)

Exemplo de configuração (adicione aos parâmetros do job):

table.optimizer.window-rewrite-enabled: true
Importante

Este parâmetro serve apenas como medida transitória de compatibilidade durante atualizações. A sintaxe obsoleta pode perder o suporte à reescrita em versões futuras. Migre para a nova sintaxe o mais breve possível.

Agregação por TVF de janela

A agregação por TVF de janela define janelas por meio de uma cláusula GROUP BY que inclui as colunas window_start e window_end produzidas pelas TVFs de janela. Ela corresponde ao operador WindowAggregate e oferece suporte às funções de janela TUMBLE, HOP, CUMULATE e SESSION.

Diferentemente da agregação em tabelas contínuas, a agregação por TVF de janela não gera resultados intermediários, apenas um resultado final ao término de cada janela. O sistema limpa automaticamente os dados de estado intermediários.

Para obter detalhes sobre sintaxe, exemplos e recursos, consulte Agregação por TVF de janela.

Sintaxe da TVF de janela SESSION: VVR 11.x vs VVR 8.x

A sintaxe da TVF de janela SESSION varia entre as versões do VVR. Atualize para o VVR 11.1 ou posterior para utilizar a sintaxe completa.

VVR 11.x (Flink 1.20)

SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap)

Parâmetro

Descrição

data

Uma tabela com uma coluna de atributo de tempo

keycols

(Opcional) Colunas usadas para particionar dados antes do janelamento de sessão

timecol

A coluna de atributo de tempo mapeada para janelas de sessão

gap

O intervalo máximo de tempo entre dois eventos na mesma sessão

VVR 8.x (Flink 1.17)

SESSION(TABLE data, DESCRIPTOR(timecol), gap)

Parâmetro

Descrição

data

Uma tabela com uma coluna de atributo de tempo

timecol

A coluna de atributo de tempo mapeada para janelas de sessão

gap

O intervalo máximo de tempo entre dois eventos na mesma sessão

O VVR 8.x não oferece suporte a PARTITION BY . Os campos de partição são inferidos implicitamente a partir da cláusula GROUP BY .

Comparação de sintaxe SESSION: VVR 11.x vs VVR 8.x

VVR 11.x

VVR 8.x

Sintaxe

SESSION(TABLE data [PARTITION BY(keycols, ...)], DESCRIPTOR(timecol), gap)

SESSION(TABLE data, DESCRIPTOR(timecol), gap)

Especificação do campo de partição

Explícita — via PARTITION BY(keycols)

Implícita — via cláusula GROUP BY

Restrições do campo de partição

Nenhuma

Deve estar em GROUP BY; não pode ser window_start, window_end ou window_time

Uso independente de SESSION()

Suportado

Deve ser usado com GROUP BY

Mesclagem de função de janela com agregação

Suportada

Não suportada — a agregação deve corresponder aos campos de partição

Os exemplos a seguir são equivalentes. Ambos usam item como campo de partição.

-- The Bid table schema (used in all examples below)
> desc Bid;
+-------------+------------------------+------+-----+--------+---------------------------------+
|        name |                   type | null | key | extras |                       watermark |
+-------------+------------------------+------+-----+--------+---------------------------------+
|     bidtime | TIMESTAMP(3) *ROWTIME* | true |     |        | `bidtime` - INTERVAL '1' SECOND |
|       price |         DECIMAL(10, 2) | true |     |        |                                 |
|        item |                 STRING | true |     |        |                                 |
+-------------+------------------------+------+-----+--------+---------------------------------+

-- VVR 11.x: partition field declared explicitly in SESSION()
> SELECT window_start, window_end, item, SUM(price) AS total_price
  FROM TABLE(
      SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
  GROUP BY item, window_start, window_end;

-- VVR 8.x: partition field inferred from GROUP BY
> SELECT window_start, window_end, item, SUM(price) AS total_price
  FROM TABLE(
      SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
  GROUP BY item, window_start, window_end;

VVR 11.x

VVR 8.x

Particionamento de janela SESSION

SESSION(TABLE Bid PARTITION BY item, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)

SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)

Agregação e mesclagem de janelas

Mesclagem direta suportada (ex.: SUM(price) dentro da janela)

Os campos de agregação devem corresponder aos campos de partição da janela (ex.: GROUP BY item)

Quando a agregação por TVF de janela reverte para o modo não-TVF

Quando uma consulta inclui uma TVF de janela, mas não atende às condições para mesclar a TVF e a agregação, o sistema reverte para um plano de execução não-TVF.

Aviso

Se uma consulta não mesclável usar processing time como atributo de tempo, a coluna de processing time será materializada e utilizada como atributo de tempo das janelas criadas. Isso faz com que a watermark da tabela de source afete os resultados da agregação: as janelas podem fechar antes do esperado e dados atrasados podem ser descartados, da mesma forma que ocorre com janelas de event-time. Evite os padrões abaixo para prevenir esse comportamento.

A TVF de janela e a instrução de agregação não podem ser mescladas quando qualquer uma das seguintes condições for atendida:

  1. Filtragem ou computação em campos de tempo da janela. window_start, window_end ou window_time é filtrado ou modificado antes da agregação.

    -- Filtering on window_start
    > SELECT window_start, window_end, item, SUM(price) AS total_price
        FROM
        (SELECT item, price, window_start, window_end FROM
        TABLE(
        SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
        WHERE window_start >= TIMESTAMP '2020-04-15 08:06:00.000')
        GROUP BY item, window_start, window_end;
    
    -- Arithmetic on window_start
    > SELECT window_start, window_end, item, SUM(price) AS total_price
        FROM
        (SELECT item, price, window_start + (INTERVAL '1' SECOND) AS window_start, window_end FROM
        TABLE(
        SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)))
        GROUP BY item, window_start, window_end;
    
    -- Type casting on window_start
    > SELECT window_start, window_end, item, SUM(price) AS total_price
        FROM
        (SELECT item, price, CAST(window_start AS varchar) AS window_start, window_end FROM
        TABLE(
        SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)))
        GROUP BY item, window_start, window_end;
  2. Uma TVF de janela é usada com uma função com valor de tabela definida pelo usuário (UDTF).

    > SELECT window_start, window_end, category, SUM(price) AS total_price
        FROM
        (SELECT category, price, window_start, window_end FROM
        TABLE(
        SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES)),
        LATERAL TABLE(category_udtf(item)) AS T(category))
        GROUP BY category, window_start, window_end;
  3. A cláusula GROUP BY não contém window_start ou window_end.

    > SELECT window_start, item, SUM(price) AS total_price
      FROM TABLE(
          SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
      GROUP BY item, window_start;
  4. Uma função de agregação definida pelo usuário (UDAF) Python é utilizada.

  5. GROUPING SETS, CUBE ou ROLLUP agrupam separadamente por window_start ou window_end.

    > SELECT item, SUM(price) AS total_price
      FROM TABLE(
          SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
      GROUP BY GROUPING SETS((item), (window_start), (window_end));
    
    > SELECT item, SUM(price) AS total_price
      FROM TABLE(
          SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
      GROUP BY CUBE (item, window_start, window_end);
    
    > SELECT item, SUM(price) AS total_price
      FROM TABLE(
          SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
      GROUP BY ROLLUP (item, window_start, window_end);
  6. Uma função de agregação é aplicada a window_start, window_end ou window_time.

    > SELECT window_start, window_end, item, SUM(price) AS total_price, MAX(window_end) AS max_end
      FROM TABLE(
          SESSION(TABLE Bid, DESCRIPTOR(bidtime), INTERVAL '5' MINUTES))
      GROUP BY item, window_start, window_end;

Suporte a stream de atualização

A tabela abaixo mostra o suporte a stream de atualização por função de janela e sintaxe.

Função de janela

Sintaxe antiga (GroupWindowAggregation) — VVR

Sintaxe antiga (GroupWindowAggregation) — Apache Flink

Nova sintaxe (WindowAggregate) — VVR

Nova sintaxe (WindowAggregate) — Apache Flink

TUMBLE

Sim

Sim

Sim

Não

HOP

Sim

Sim

Sim

Não

SESSION

Sim

Sim

Sim

Sim (Apache Flink 1.19 e posterior)

CUMULATE

N/A

N/A

Sim (VVR 8.0.6 e posterior)

Não

Na sintaxe antiga, o suporte a stream de atualização é idêntico tanto no VVR quanto no Apache Flink. Na nova sintaxe, apenas o operador WindowAggregate do VVR oferece suporte a streams de atualização para todas as funções de janela. O VVR selecione automaticamente entre os operadores GroupWindowAggregation e WindowAggregate com base na stream de entrada.

Para diferenças entre a função de janela SESSION no VVR e no Apache Flink, consulte Consultas .

Otimização Local-Global para agregação de janela

A agregação de janela suporta a otimização de agregação em duas fases Local-Global. Quando habilitada, o otimizador divide uma agregação de janela de fase única em:

  • Local Aggregate: realiza uma pré-agregação parcial antes do shuffle de dados, reduzindo a quantidade de dados transferidos pela rede.

  • Global Aggregate: executa a agregação final após o shuffle e produz os resultados.

Todas as seis condições a seguir devem ser atendidas para que essa otimização tenha efeito.

Condição 1: A estratégia de fase de agregação permite duas fases

table.optimizer.agg-phase-strategy está definido como AUTO (padrão) ou TWO_PHASE. Defina como ONE_PHASE desabilita a otimização de duas fases.

Condição 2: A janela usa event time

A janela deve usar event time (rowtime). Janelas de processing-time não são suportadas.

Condição 3: O tipo de janela não é SESSION

Os tipos de janela TUMBLE, HOP e CUMULATE são suportados. Janelas SESSION não oferecem suporte à otimização de duas fases.

Condição 4: Todas as funções de agregação suportam mesclagem parcial

Todas as funções de agregação devem suportar a operação de mesclagem. Funções integradas como SUM, COUNT, MIN, MAX e AVG são suportadas. UDAFs personalizadas devem implementar o método merge().

Condição 5: A stream de entrada é insert-only e a janela pode ser convertida para a forma TVF

As seguintes condições devem ser atendidas simultaneamente:

  • A stream de entrada é insert-only.

  • table.exec.emit.early-fire.enabled está definido como false (padrão).

  • table.exec.emit.late-fire.enabled está definido como false (padrão).

  • Para janelas HOP, a janela deve estar alinhada (o tamanho da janela deve ser divisível pelo intervalo de deslizamento).

Condição 6: A distribuição de dados ainda não satisfaz os requisitos de particionamento

A distribuição dos dados de entrada ainda não atende aos requisitos de particionamento para agregação. Se os dados já estiverem distribuídos pela chave de partição, o otimizador determina que nenhuma pré-agregação adicional é necessária e não gera um nó Local Aggregate.

Referência de parâmetros

Parâmetro

Tipo

Padrão

Requisito

table.optimizer.agg-phase-strategy

Enum

AUTO

Não deve ser ONE_PHASE

table.exec.emit.early-fire.enabled

Boolean

false

Deve ser false

table.exec.emit.late-fire.enabled

Boolean

false

Deve ser false