Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Python

Última atualização: Jun 27, 2026

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")
Nota

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 python.fn-execution.bundle.size define o número máximo de registros armazenáveis em cache.

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 python.fn-execution.bundle.size ou quando o tempo de cache atinge o limiar de python.fn-execution.bundle.time.

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 python.fn-execution.arrow.batch.size não pode ser maior que o valor do parâmetro python.fn-execution.bundle.size.

Nota

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