Este tópico descreve como criar, registrar e usar uma função de valor de tabela definida pelo usuário (UDTF) Python no Realtime Compute for Apache Flink.
Descrição
Uma UDTF aceita zero, um ou mais valores escalares como parâmetros de entrada, que podem ter comprimento variável. Diferentemente das funções escalares definidas pelo usuário (UDFs), as UDTFs retornam qualquer número de linhas em vez de um único valor. As linhas retornadas podem incluir uma ou mais colunas, resultando em múltiplas linhas ou colunas a cada chamada da função.
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 cada workspace. VVR anterior à 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 UDTF
O Flink fornece códigos de exemplo de extensões definidas pelo usuário (UDXs) Python para auxiliar no desenvolvimento de UDXs. O código de exemplo inclui implementações de UDFs Python, funções de agregação definidas pelo usuário (UDAFs) Python e UDTFs Python. Esta seção descreve como criar uma UDTF no sistema operacional Windows.
Baixe e descompacte o arquivo python_demo-master na sua máquina local.
-
Clique duas vezes no arquivo udtfs.py no diretório \python_demo-master\udx e modifique o conteúdo conforme suas necessidades de negócio.
Neste exemplo, split define o código capaz de separar uma linha de string em várias colunas usando barras verticais (|).
from pyflink.table import DataTypes from pyflink.table.udf import udtf @udtf(result_types=[DataTypes.STRING(), DataTypes.STRING()]) def split(s: str): splits = s.split("|") yield splits[0], splits[1] -
Acesse o diretório \python_demo onde a pasta udx está localizada e execute o comando abaixo para compactar os arquivos do diretório:
zip -r python_demo.zip udxSe o pacote python_demo.zip aparecer no diretório \python_demo\, a UDTF foi desenvolvida com sucesso.
Registrar uma UDTF
Para obter mais informações sobre como registrar uma UDTF, consulte Gerenciar UDFs.
Usar uma UDTF
Após registrar uma UDTF, siga as etapas abaixo para utilizá-la:
-
Use Flink SQL para criar um rascunho. Para mais detalhes, consulte Visão geral do desenvolvimento de jobs.
Ao concatenar a string "aa" e o campo message de cada linha na tabela ASI_UDTF_Source com barras verticais (|), as strings resultantes são divididas em várias colunas por essas mesmas barras. O código abaixo mostra um exemplo:
CREATE TEMPORARY TABLE ASI_UDTF_Source ( `message` VARCHAR ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE ASI_UDTF_Sink ( name VARCHAR, place VARCHAR ) WITH ( 'connector' = 'blackhole' ); INSERT INTO ASI_UDTF_Sink SELECT name,place FROM ASI_UDTF_Source,lateral table(split(concat_ws('|', `message`, 'aa'))) as T(name,place); -
No painel de navegação à esquerda do console de desenvolvimento do Realtime Compute for Apache Flink, escolha . Na página Deployments, localize o deployment desejado e clique em Start na coluna Actions.
Depois que o deployment for iniciado, duas colunas de dados serão inseridas na tabela ASI_UDTF_Sink. Essas colunas contêm as strings concatenadas separadas por barras verticais (|).