Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:UDAFs em Python

Última atualização: Jun 27, 2026

Uma função de agregação definida pelo usuário (UDAF) agrega vários valores de entrada em um único valor de saída, mapeando diversas linhas de entrada para um resultado por grupo.

Este tópico explica como criar, registrar e usar uma UDAF em Python no Realtime Compute for Apache Flink.

Limitações

  • É necessário ter o Apache Flink 1.12 ou posterior.

  • O Python já vem pré-instalado no workspace do Realtime Compute for Apache Flink. Escreva seu código na versão do Python pré-instalada.

    O Python 3.7.9 é pré-instalado nas versões do Ververica Runtime (VVR) anteriores à 8.0.11. O Python 3.9.21 é pré-instalado no VVR 8.0.11 ou posterior. Após atualizar para o VVR 8.0.11 ou superior, teste novamente, reimplante e execute novamente todos os jobs PyFlink criados em uma versão anterior do VVR.
  • Os ambientes de execução suportam JDK 8 e JDK 11. Caso sua implantação em Python dependa de um arquivo JAR de terceiros, verifique se esse arquivo é compatível com o JDK 8 ou JDK 11.

  • Apenas a versão open source Scala 2.11 é suportada. Se a sua implantação em Python depender de um arquivo JAR de terceiros, certifique-se de que ele seja compatível com o Scala 2.11.

Funcionamento

Uma UDAF utiliza um acumulador para rastrear o estado intermediário da agregação ao longo das linhas de entrada. O acumulador é criado uma vez por grupo e atualizado linha a linha até que todas as entradas sejam processadas.

A ordem de execução para cada grupo é:

  1. create_accumulator() — cria um acumulador vazio para armazenar o estado inicial.

  2. accumulate(accumulator, ...) — chamado uma vez para cada linha de entrada, a fim de atualizar o acumulador.

  3. get_value(accumulator) — invocado após o processamento de todas as linhas do grupo para retornar o resultado final.

Referência de métodos

Método

Obrigatório

Quando implementar

create_accumulator()

Sim

Sempre

accumulate(...)

Sim

Sempre

get_value(...)

Sim

Sempre

retract(...)

Condicional

Quando a agregação puder receber mensagens de retração.

Crie uma UDAF

O Flink fornece códigos de exemplo para extensões definidas pelo usuário (UDXs), incluindo UDAFs, funções definidas pelo usuário (UDFs) e funções de tabela definidas pelo usuário (UDTFs). Os passos abaixo utilizam um ambiente Windows.
  1. Baixe e descompacte o pacote python_demo-master em sua máquina.

  2. No PyCharm, escolha File > Open e abra o diretório descompactado python_demo-master.

  3. Abra o arquivo \python_demo-master\udx\udfs.py e edite-o conforme sua lógica de negócios. O exemplo abaixo define weighted_avg, que calcula uma média ponderada sobre dados atuais e históricos.

    from pyflink.common import Row
    from pyflink.table import AggregateFunction, DataTypes
    from pyflink.table.udf import udaf
    
    class WeightedAvg(AggregateFunction):
    
        def create_accumulator(self):
            # Row(sum, count)
            return Row(0, 0)
    
        def get_value(self, accumulator: Row) -> float:
            if accumulator[1] == 0:
                return 0
            else:
                return accumulator[0] / accumulator[1]
    
        def accumulate(self, accumulator: Row, value, weight):
            accumulator[0] += value * weight
            accumulator[1] += weight
    
        def retract(self, accumulator: Row, value, weight):
            accumulator[0] -= value * weight
            accumulator[1] -= weight
    
    weighted_avg = udaf(f=WeightedAvg(),
                        result_type=DataTypes.DOUBLE(),
                        accumulator_type=DataTypes.ROW([
                            DataTypes.FIELD("f0", DataTypes.BIGINT()),
                            DataTypes.FIELD("f1", DataTypes.BIGINT())]))
  4. No diretório \python_demo-master\, empacote a pasta udx:

    zip -r python_demo.zip udx

    A UDAF será criada quando o arquivo python_demo.zip aparecer no diretório \python_demo-master\.

Registrar uma UDAF

Para obter as etapas de registro, consulte Gerencie funções definidas pelo usuário (UDFs).

Usar uma UDAF

  1. Crie um rascunho de Flink SQL. Para mais detalhes, veja Visão geral do desenvolvimento de jobs. O exemplo a seguir calcula a média ponderada do campo a em ASI_UDAF_Source, utilizando o campo b como peso.

    CREATE TEMPORARY TABLE ASI_UDAF_Source (
      a   BIGINT,
      b   BIGINT
    ) WITH (
      'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE ASI_UDAF_Sink (
      avg_value  DOUBLE
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO ASI_UDAF_Sink
    SELECT weighted_avg(a, b)
    FROM ASI_UDAF_Source;
  2. No painel de navegação à esquerda do console de desenvolvimento do Realtime Compute for Apache Flink, selecione O&M > Deployments. Localize a implantação desejada e clique em Start na coluna Actions. Após o início da implantação, a média ponderada do campo a — tendo o campo b como peso — será gravada em cada linha de ASI_UDAF_Sink.