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 é:
create_accumulator()— cria um acumulador vazio para armazenar o estado inicial.accumulate(accumulator, ...)— chamado uma vez para cada linha de entrada, a fim de atualizar o acumulador.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 |
|
|
Sim |
Sempre |
|
|
Sim |
Sempre |
|
|
Sim |
Sempre |
|
|
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.
Baixe e descompacte o pacote python_demo-master em sua máquina.
No PyCharm, escolha File > Open e abra o diretório descompactado
python_demo-master.-
Abra o arquivo
\python_demo-master\udx\udfs.pye edite-o conforme sua lógica de negócios. O exemplo abaixo defineweighted_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())])) -
No diretório
\python_demo-master\, empacote a pastaudx:zip -r python_demo.zip udxA UDAF será criada quando o arquivo
python_demo.zipaparecer 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
-
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
aemASI_UDAF_Source, utilizando o campobcomo 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; 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 campobcomo peso — será gravada em cada linha deASI_UDAF_Sink.