O Realtime Compute for Apache Flink oferece suporte a funções definidas pelo usuário (UDFs) em Python nos jobs do Flink SQL. Crie funções escalares, de agregação e de tabela em Python, gerencie dependências do Python e ajuste o desempenho das UDFs.
Tipos de funções definidas pelo usuário
|
Categoria |
Descrição |
|
função escalar definida pelo usuário (UDSF) |
Uma UDSF mapeia zero, um ou mais valores escalares para um novo valor escalar. Processa uma linha de entrada e produz um valor de saída, estabelecendo um mapeamento de um para um. Para obter mais informações, consulte Funções escalares definidas pelo usuário (UDSFs). |
|
função de agregação definida pelo usuário (UDAF) |
Uma UDAF agrega vários registros em um único registro, criando um mapeamento de muitos para um. Para saber mais, consulte Funções de agregação definidas pelo usuário (UDAFs). |
|
função de tabela definida pelo usuário (UDTF) |
Uma UDTF aceita zero, um ou mais valores escalares como parâmetros de entrada. Diferentemente de uma função escalar, pode retornar qualquer número de linhas, cada uma composta por uma ou mais colunas. Para obter detalhes, consulte Funções de tabela definidas pelo usuário (UDTFs). |
Uso de dependências do Python
Os clusters do Realtime Compute for Apache Flink já incluem pacotes Python comuns pré-instalados, como Pandas, NumPy e PyArrow. Consulte Desenvolver jobs em Python para ver a lista completa de pacotes Python de terceiros disponíveis. Importe esses pacotes na função antes de usá-los, conforme o exemplo a seguir.
@udf(result_type=DataTypes.FLOAT())
def percentile(values: List[float], percentile: float):
import numpy as np
return np.percentile(values, percentile)
Para usar um pacote Python de terceiros não pré-instalado, faça upload dele como arquivo de dependência ao registrar a UDF em Python. Para mais detalhes, consulte Gerenciar funções definidas pelo usuário (UDFs) e Usar dependências do Python.
Depuração de código
Use o módulo logging para gerar informações de log da função definida pelo usuário em Python e facilitar a solução de problemas. O exemplo abaixo demonstra como aplicar o logging.
@udf(result_type=DataTypes.BIGINT())
def add(i, j):
logging.info("hello world")
return i + j
Após a geração dos logs, visualize-os nos arquivos de log do TaskManager. Para obter instruções detalhadas, consulte Visualizar logs de execução.
Ajuste de desempenho
Pré-carregamento de recursos
Carregue os recursos durante a inicialização da função para evitar recarregamentos a cada chamada do método eval. Essa abordagem permite, por exemplo, carregar um modelo grande de deep learning apenas uma vez e executar previsões em lote posteriormente.
from pyflink.table import DataTypes
from pyflink.table.udf import ScalarFunction, udf
class Predict(ScalarFunction):
def open(self, function_context):
import pickle
with open("resources.zip/resources/model.pkl", "rb") as f:
self.model = pickle.load(f)
def eval(self, x):
return self.model.predict(x)
predict = udf(Predict(), result_type=DataTypes.DOUBLE(), func_type="pandas")
Para saber como fazer upload de arquivos de dados em Python, consulte Usar dependências do Python.
Funções definidas pelo usuário assíncronas
Em cenários com uso intensivo de I/O, como acesso a bancos de dados externos ou chamadas a serviços HTTP, use funções definidas pelo usuário assíncronas. Uma única instância da função processa múltiplas requisições simultaneamente, distribuindo o tempo de espera entre várias chamadas e aumentando significativamente o throughput do job. Esse recurso está disponível apenas no VVR 11.7 e versões posteriores, aplicando-se exclusivamente a funções escalares definidas pelo usuário (UDSFs). Para obter detalhes, consulte Funções definidas pelo usuário assíncronas.
Uso da biblioteca Pandas
Além das funções definidas pelo usuário padrão em Python, o Realtime Compute for Apache Flink oferece suporte a UDFs baseadas em Pandas. Essas funções aceitam estruturas de dados do Pandas, como pandas.Series e pandas.DataFrame, como entrada, permitindo o uso de bibliotecas de alto desempenho como Pandas e NumPy. Para obter mais informações, consulte Funções Definidas pelo Usuário Vetorizadas.
Parâmetros
O desempenho de uma função definida pelo usuário em Python depende principalmente da implementação. Em caso de problemas de performance, otimize primeiro a lógica da função. Os parâmetros listados abaixo também influenciam o desempenho.
|
Parâmetro |
Descrição |
|
python.fn-execution.bundle.size |
As UDFs em Python executam de forma assíncrona. O operador Java armazena dados em cache antes de enviá-los a um processo Python para execução. Quando o cache atinge um determinado limiar, os dados são enviados ao processo Python. O parâmetro O valor padrão é 100.000 registros. |
|
python.fn-execution.bundle.time |
Este parâmetro controla o tempo máximo de armazenamento em cache. Os dados armazenados são enviados para processamento quando a contagem de registros atinge o limiar de O valor padrão é 1.000 milissegundos. |
|
python.fn-execution.arrow.batch.size |
Para UDFs com Pandas, este parâmetro especifica o número máximo de registros que um lote Arrow pode conter. O valor padrão é 10.000. Nota
O valor do parâmetro |
Definir esses parâmetros com valores excessivamente altos pode ser contraproducente. O buffering excessivo durante um checkpoint pode causar demora ou falha nos checkpoints. Para mais informações sobre esses parâmetros, consulte Configuração.
Tópicos relacionados
Para registrar, atualize e exclua uma função definida pelo usuário, consulte Gerenciar funções definidas pelo usuário (UDFs).
Para demonstrações sobre desenvolvimento e uso de funções definidas pelo usuário em Python, consulte Funções de agregação definidas pelo usuário (UDAFs), Funções escalares definidas pelo usuário (UDSFs) e Funções de tabela definidas pelo usuário (UDTFs).
Para aprender a usar ambientes virtuais Python personalizados, pacotes Python de terceiros, pacotes JAR e arquivos de dados em um job Flink Python, consulte Usar dependências do Python.
Para exemplos de desenvolvimento e uso de funções definidas pelo usuário em Java, consulte Funções de agregação definidas pelo usuário (UDAFs), Funções escalares definidas pelo usuário (UDSFs) e Funções de tabela definidas pelo usuário (UDTFs).
Para depurar e ajustar funções definidas pelo usuário em Java, consulte Visão geral.