Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Funções escalares definidas pelo usuário (UDSFs) em Python

Última atualização: Jun 27, 2026

Uma função escalar definida pelo usuário (UDSF) mapeia zero, um ou mais valores escalares para um único valor escalar. Cada linha de entrada gera exatamente um valor de saída.

Este tópico descreve como criar, registrar e usar uma UDSF Python no Realtime Compute for Apache Flink.

Limites

As seguintes restrições se aplicam ao desenvolvimento de funções definidas pelo usuário (UDFs) Python no Realtime Compute for Apache Flink:

Restrição

Requisito

Versão do Apache Flink

1.12 e posteriores

Versão do Python

Pré-instalado em todos os workspaces. VVR anterior a 8.0.11: Python 3.7.9. VVR 8.0.11 e posterior: Python 3.9.21.

Versão do JDK

JDK 8 e JDK 11. Pacotes JAR de terceiros devem ser compatíveis com JDK 8 ou JDK 11.

Versão do Scala

Apenas Scala 2.11 open-source. Pacotes JAR de terceiros devem ser compatíveis com Scala 2.11.

Importante

Após atualizar para o VVR 8.0.11 ou posterior, teste, implante e execute novamente seus rascunhos PyFlink existentes para confirmar a compatibilidade.

Criar uma UDSF

Os passos a seguir utilizam o Windows como ambiente de exemplo. O Flink fornece um repositório de amostra que inclui implementações para UDSFs, funções de agregação definidas pelo usuário (UDAFs) e funções com valor de tabela definidas pelo usuário (UDTFs).
  1. Baixe e descompacte o python_demo-master em sua máquina local.

    Este é um repositório GitHub de terceiros. O acesso pode ser lento ou intermitente.
  2. No PyCharm, escolha File > Open e abra o diretório descompactado python_demo-master.

  3. Abra o arquivo udfs.py no caminho \python_demo-master\udx e defina sua UDSF.

    from pyflink.table import DataTypes
    from pyflink.table.udf import udf
    
    @udf(result_type=DataTypes.STRING())
    def sub_string(s: str, begin: int, end: int):
        return s[begin:end]

    O exemplo sub_string extrai caracteres da posição begin até a posição end na string de entrada.

  4. No diretório \python_demo-master, execute o comando a seguir para empacotar o diretório udx:

    zip -r python_demo.zip udx

    Quando o arquivo python_demo.zip aparecer em \python_demo-master\, o pacote estará pronto.

Registrar uma UDSF

Após criar o pacote, registre a UDSF no console do Realtime Compute for Apache Flink. Para obter as etapas de registro, consulte Gerenciar funções definidas pelo usuário (UDFs).

Usar uma UDSF

Depois de registrar a UDSF, utilize-a em um job Flink SQL.

  1. Crie um rascunho usando Flink SQL. Para mais detalhes, consulte Visão geral do desenvolvimento de jobs. O exemplo a seguir chama ASI_UDSF (o nome registrado da sua UDSF) para extrair caracteres das posições 2 a 4 do campo a na tabela de origem:

    CREATE TEMPORARY TABLE ASI_UDSF_Source (
      a VARCHAR,
      b INT,
      c INT
    ) WITH (
      'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE ASI_UDSF_Sink (
      a VARCHAR
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO ASI_UDSF_Sink
    SELECT ASI_UDSF(a, 2, 4)
    FROM ASI_UDSF_Source;
  2. No painel de navegação à esquerda do console de desenvolvimento, escolha O&M > Deployments. Localize o deployment e clique em Start na coluna Actions. Após o início do deployment, os caracteres nas posições 2–4 do campo a em ASI_UDSF_Source serão gravados em ASI_UDSF_Sink.

Funções definidas pelo usuário assíncronas

Para UDFs que realizam operações intensivas em I/O, como acesso a bancos de dados externos ou requisições HTTP, utilize funções definidas pelo usuário assíncronas. Uma única UDSF assíncrona consegue processar múltiplas requisições de I/O simultaneamente, distribuindo o tempo de espera entre as requisições e melhorando o throughput do job.

Limites

  • Suportado apenas no VVR 11.7 e versões posteriores. É necessário o VVR PyFlink 11.7 ou superior. Para mais detalhes, consulte ververica-flink.

    pip3 install "ververica-flink>=11.7"

  • Apenas funções escalares definidas pelo usuário (UDSFs) assíncronas são suportadas.

  • Somente o modo de processo Python é suportado, ou seja, python.execution-mode=process.

  • Funções definidas pelo usuário assíncronas com Pandas ainda não são suportadas.

Uso

É possível implementar uma função definida pelo usuário assíncrona como uma função async do Python ou como uma subclasse da classe de função assíncrona. O código de exemplo é mostrado abaixo.

import asyncio

from pyflink.table import DataTypes
from pyflink.table.udf import AsyncScalarFunction, udf

# Method 1: Use a Python async function
@udf(result_type=DataTypes.STRING())
async def async_api_call(product_id: str) -> str:
    await asyncio.sleep(0.05)
    return f"product_{product_id}"

# Method 2: Subclass the asynchronous function class
class AsyncUserLookup(AsyncScalarFunction):
    def open(self, function_context):
        self.cache = {}

    async def eval(self, user_id: str) -> str:
        if user_id in self.cache:
            return self.cache[user_id]

        await asyncio.sleep(0.05)
        result = f"user_{user_id}"
        self.cache[user_id] = result
        return result

    def close(self):
        self.cache.clear()

async_user_lookup = udf(
    AsyncUserLookup(),
    input_types=[DataTypes.STRING()],
    result_type=DataTypes.STRING()
)

O registro e o uso de funções definidas pelo usuário assíncronas seguem o mesmo padrão das funções síncronas.

Parâmetros de configuração

Os parâmetros a seguir controlam o comportamento em tempo de execução das funções definidas pelo usuário assíncronas.

Parâmetro

Valor padrão

Descrição

table.exec.async-scalar.max-concurrent-operations

10

Número máximo de chamadas assíncronas simultâneas por instância de operador. Valor padrão: 10.

table.exec.async-scalar.timeout

3 min

Tempo limite para uma única chamada assíncrona.

table.exec.async-scalar.retry-strategy

FIXED_DELAY

Estratégia de nova tentativa após falha em uma chamada assíncrona. Valores suportados:

  • FIXED_DELAY: Tenta novamente após um tempo de espera fixo.

  • NO_RETRY: Não tenta novamente.

table.exec.async-scalar.retry-delay

100 ms

Tempo de espera para novas tentativas com atraso fixo.

Nota

Efetivo apenas quando table.exec.async-scalar.retry-strategy está definido como FIXED_DELAY.

table.exec.async-scalar.max-attempts

3

Número máximo de tentativas antes que uma chamada assíncrona seja considerada falha.

Nota

Efetivo apenas quando table.exec.async-scalar.retry-strategy está definido como FIXED_DELAY.