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 |
|
|
|
Funções de janela |
TUMBLE, HOP, SESSION |
TUMBLE, HOP, CUMULATE, SESSION |
|
Status |
Obsoleto |
Recomendado |
|
Otimizações de desempenho |
Não |
Sim |
|
Suporte a |
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.
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 |
|
|
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
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 |
|
|
Uma tabela com uma coluna de atributo de tempo |
|
|
(Opcional) Colunas usadas para particionar dados antes do janelamento de sessão |
|
|
A coluna de atributo de tempo mapeada para janelas de sessão |
|
|
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 |
|
|
Uma tabela com uma coluna de atributo de tempo |
|
|
A coluna de atributo de tempo mapeada para janelas de sessão |
|
|
O intervalo máximo de tempo entre dois eventos na mesma sessão |
O VVR 8.x não oferece suporte aPARTITION BY. Os campos de partição são inferidos implicitamente a partir da cláusulaGROUP BY.
Comparação de sintaxe SESSION: VVR 11.x vs VVR 8.x
|
VVR 11.x |
VVR 8.x |
|
|
Sintaxe |
|
|
|
Especificação do campo de partição |
Explícita — via |
Implícita — via cláusula |
|
Restrições do campo de partição |
Nenhuma |
Deve estar em |
|
Uso independente de |
Suportado |
Deve ser usado com |
|
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 |
|
|
|
Agregação e mesclagem de janelas |
Mesclagem direta suportada (ex.: |
Os campos de agregação devem corresponder aos campos de partição da janela (ex.: |
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.
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:
-
Filtragem ou computação em campos de tempo da janela.
window_start,window_endouwindow_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; -
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; -
A cláusula
GROUP BYnão contémwindow_startouwindow_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; Uma função de agregação definida pelo usuário (UDAF) Python é utilizada.
-
GROUPING SETS,CUBEouROLLUPagrupam separadamente porwindow_startouwindow_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); -
Uma função de agregação é aplicada a
window_start,window_endouwindow_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.enabledestá definido comofalse(padrão).table.exec.emit.late-fire.enabledestá definido comofalse(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 |
|
|
Enum |
AUTO |
Não deve ser ONE_PHASE |
|
|
Boolean |
false |
Deve ser false |
|
|
Boolean |
false |
Deve ser false |