Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:UDAFs Java

Última atualização: Jun 27, 2026

Uma função de agregação definida pelo usuário (UDAF) reduz várias linhas de entrada a um único valor de saída, realizando um mapeamento de muitos para um. Utilize UDAFs quando as funções de agregação integradas, como SUM ou MAX, não atenderem à sua lógica de agregação.

Para obter informações gerais sobre funções definidas pelo usuário no Flink, consulte Funções definidas pelo usuário.

Funções definidas pelo usuário e ASI_UDX_Demo estão hospedados em sites de terceiros e podem apresentar lentidão ou indisponibilidade ocasional.

Funcionamento

Uma UDAF utiliza um acumulador para armazenar o estado intermediário da agregação. Para cada grupo de linhas que compartilha a mesma chave GROUP BY, o runtime executa três métodos em sequência:

  1. createAccumulator() — cria um novo acumulador com o estado inicial.

  2. accumulate(acc, ...) — chamado uma vez por linha de entrada para atualizar o acumulador.

  3. getValue(acc) — invocado após o processamento de todas as linhas para retornar o resultado final.

Por exemplo, se uma tabela possuir uma coluna numérica e você desejar uma soma cumulativa para os valores 1, 2 e 3 no mesmo grupo:

  • O runtime chama createAccumulator() uma única vez para inicializar o acumulador com sum = 0.

  • Em seguida, executa accumulate() para cada linha, atualizando sum para 1, depois 3 e, por fim, 6.

  • Por último, invoca getValue() para devolver o resultado final.

A saída varia conforme a ativação do mini-batch:

  • Sem mini-batch (padrão): gera uma saída por linha — 1, 3, 6.

  • Com mini-batch ativado: emite apenas o resultado final — 6. A quantidade de saídas intermediárias depende das configurações de mini-batch e da distribuição dos dados.

Para detalhes sobre a configuração de mini-batch, consulte Otimizar Flink SQL.

Métodos obrigatórios

Toda implementação de AggregateFunction deve definir estes três métodos:

Método

Finalidade

createAccumulator()

Retorna um novo acumulador com o estado inicial

accumulate(acc, ...)

Atualiza o acumulador com uma linha de entrada

getValue(acc)

Devolve o resultado agregado final

Outros métodos podem ser implementados conforme a necessidade do seu caso de uso:

Método

Finalidade

retract(acc, ...)

Permite retrair uma mensagem gerada por um operador upstream

merge(acc, iterable)

Suporta a otimização de agregação em dois estágios (local-global)

Crie uma UDAF

O Realtime Compute for Apache Flink disponibiliza um projeto de demonstração de UDF (ASI_UDX_Demo) com ambiente de desenvolvimento pré-configurado, eliminando a necessidade de configurar o ambiente manualmente.

O projeto de demonstração inclui implementações de exemplo para funções escalares definidas pelo usuário (UDSFs), UDAFs e funções de tabela definidas pelo usuário (UDTFs).

Pré-requisitos

Antes de começar, certifique-se de ter:

  • IntelliJ IDEA instalado

  • Maven instalado

  • Ambiente de desenvolvimento Java configurado

Etapas

  1. Baixe e descompacte o ASI_UDX_Demo em sua máquina local. A pasta descompactada ASI_UDX-main contém:

    • pom.xml — configuração do projeto Maven, incluindo coordenadas, dependências e regras de build.

    • \ASI_UDX-main\src\main\java\ASI_UDAF\ASI_UDAF.java — implementação de exemplo de UDAF.

  2. No IntelliJ IDEA, selecione File > Open e escolha a pasta ASI_UDX-main.

  3. Abra o arquivo pom.xml no diretório \ASI_UDX-main\ e configure as dependências. O arquivo já inclui a dependência mínima para o Flink 1.12:

    • Caso seu job não exija dependências adicionais, pule para a próxima etapa.

    • Se forem necessárias mais dependências, adicione-as ao pom.xml.

    <dependencies>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-common</artifactId>
            <version>1.12.7</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>

    Utilize a versão secundária mais recente da versão principal do Apache Flink correspondente à sua versão do Ververica Runtime (VVR). Para verificar o mapeamento entre versões do VVR e do Flink, consulte Visão geral.

  4. Abra \ASI_UDX-main\src\main\java\ASI_UDAF\ASI_UDAF.java e implemente sua lógica de agregação. O exemplo abaixo realiza uma soma cumulativa:

    package ASI_UDAF;
    
    import org.apache.flink.table.functions.AggregateFunction;
    
    import java.util.Iterator;
    
    public class ASI_UDAF{
        public static class AccSum{
            public long sum;
        }
    
        public static class MySum extends AggregateFunction<Long, AccSum>{
    
            @Override
            public Long getValue(AccSum acSum){
                return acSum.sum;
            }
    
            @Override
            public AccSum createAccumulator(){
                AccSum acCount= new AccSum();
                acCount.sum=0;
                return acCount;
            }
    
            public void accumulate(AccSum acc,long num){
                acc.sum += num;
            }
    
            /**
            * Supports retracting a message generated by an upstream operator.
            */
            public void retract(AccSum acc,long num){
                acc.sum -= num;
            }
    
            /**
            * Supports local-global two-stage aggregate optimization.
            */
            public void merge(AccSum acc,Iterable<AccSum> it){
                Iterator<AccSum> iter=it.iterator();
                while(iter.hasNext()){
                    AccSum accSum=iter.next();
                    if(null!=accSum){
                        acc.sum+=accSum.sum;
                    }
                }
            }
        }
    }
  5. No diretório onde está localizado o pom.xml, execute:

    mvn package -Dcheckstyle.skip

    O empacotamento da UDAF será bem-sucedido quando o arquivo ASI_UDX-1.0-SNAPSHOT.jar aparecer em \ASI_UDX-main\target\.

Usar uma UDAF

Existem duas formas de utilizar uma UDAF em deployments SQL. A tabela a seguir resume as principais diferenças:

Aspecto

Método 1: UDAF registrada

Método 2: JAR no nível do deployment

Escopo

Disponível em múltiplos deployments

Apenas um único deployment

Como registrar

Registre pela página de gerenciamento de UDFs

Faça upload do JAR em Additional dependency files no deployment

Como chamar no SQL

Chame pelo nome registrado (sem necessidade de CREATE TEMPORARY FUNCTION)

Defina um alias com CREATE TEMPORARY FUNCTION … AS 'fully.qualified.ClassName'

Reutilização

Alta — ideal para lógicas de negócio compartilhadas

Baixa — vinculada a um único deployment

Método 1: Usar uma UDAF registrada (recomendado)

Registre a UDAF uma única vez e reutilize-a em vários deployments. Para instruções de registro, consulte Gerencie UDFs.

Após registrá-la como ASI_UDAF$MySum, chame-a diretamente no seu SQL:

CREATE TEMPORARY TABLE ASI_UDAF_Source (
  a BIGINT NOT NULL
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE ASI_UDAF_Sink (
  sum  BIGINT
) WITH (
  'connector' = 'print'
);

INSERT INTO ASI_UDAF_Sink
SELECT `ASI_UDAF$MySum`(a)
FROM ASI_UDAF_Source;

Método 2: Fazer upload de um JAR para um deployment específico

Na página Data Studio > ETL do Flink, faça upload do pacote JAR usando Additional dependency files em More configurations. Em seguida, defina um alias de função temporária no SQL do job.

Esse JAR fica restrito exclusivamente àquele deployment, não podendo ser compartilhado com outros.

Se a função temporária for nomeada como mysum:

CREATE TEMPORARY TABLE ASI_UDAF_Source (
  a   BIGINT
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE ASI_UDAF_Sink (
  sum  BIGINT
) WITH (
  'connector' = 'print'
);

CREATE TEMPORARY FUNCTION `mysum` AS 'ASI_UDAF.ASI_UDAF$MySum'; -- Create the temporary function mysum.

INSERT INTO ASI_UDAF_Sink
SELECT `mysum`(a)
FROM ASI_UDAF_Source;

Execute o job

Após desenvolver e implantar o job SQL, acesse Operation Center > Job O&M. Localize o job desejado e clique em Start na coluna Operation.

Depois que o job for iniciado, a UDAF agregará o campo a de ASI_UDAF_Source e gravará a soma cumulativa em ASI_UDAF_Sink.

Próximos passos

  • Gerencie UDFs — registre e gerencie UDAFs para reutilização em diferentes deployments.

  • Otimizar Flink SQL — configure o mini-batch para controlar a emissão de resultados intermediários.

  • Visão geral — encontre a versão do Apache Flink compatível com sua versão do VVR.