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. |
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).
-
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.
No PyCharm, escolha File > Open e abra o diretório descompactado
python_demo-master.-
Abra o arquivo
udfs.pyno caminho\python_demo-master\udxe 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_stringextrai caracteres da posiçãobeginaté a posiçãoendna string de entrada. -
No diretório
\python_demo-master, execute o comando a seguir para empacotar o diretórioudx:zip -r python_demo.zip udxQuando o arquivo
python_demo.zipaparecer 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.
-
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 campoana 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; 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
aemASI_UDSF_Sourceserão gravados emASI_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:
|
|
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. |