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:
createAccumulator()— cria um novo acumulador com o estado inicial.accumulate(acc, ...)— chamado uma vez por linha de entrada para atualizar o acumulador.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 comsum = 0.Em seguida, executa
accumulate()para cada linha, atualizandosumpara1, depois3e, 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 |
|
|
Retorna um novo acumulador com o estado inicial |
|
|
Atualiza o acumulador com uma linha de entrada |
|
|
Devolve o resultado agregado final |
Outros métodos podem ser implementados conforme a necessidade do seu caso de uso:
|
Método |
Finalidade |
|
|
Permite retrair uma mensagem gerada por um operador upstream |
|
|
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
-
Baixe e descompacte o ASI_UDX_Demo em sua máquina local. A pasta descompactada
ASI_UDX-mainconté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.
No IntelliJ IDEA, selecione File > Open e escolha a pasta
ASI_UDX-main.-
Abra o arquivo
pom.xmlno 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.
-
Abra
\ASI_UDX-main\src\main\java\ASI_UDAF\ASI_UDAF.javae 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; } } } } } -
No diretório onde está localizado o
pom.xml, execute:mvn package -Dcheckstyle.skipO empacotamento da UDAF será bem-sucedido quando o arquivo
ASI_UDX-1.0-SNAPSHOT.jaraparecer 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 |
Defina um alias com |
|
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.