Uma função escalar definida pelo usuário (UDSF) mapeia zero, um ou mais valores escalares para um novo valor escalar, estabelecendo uma relação de um para um entre as linhas de entrada e os valores de saída. Use uma UDSF quando as funções internas do Flink SQL não conseguirem expressar sua lógica personalizada.
Para obter uma visão geral de todos os tipos de funções definidas pelo usuário (UDF), consulte Funções definidas pelo usuário.
Como funciona
Uma UDSF estende a classe ScalarFunction do Apache Flink e implementa um ou mais métodos eval(). O Flink chama o método eval() uma vez para cada linha de entrada e usa o valor retornado como o escalar de saída.
public class SubstringFunction extends ScalarFunction {
public String eval(String s, Integer begin, Integer end) {
return s.substring(begin, end);
}
}
Requisitos do método eval():
Declarado como
publicSuporta sobrecarga de método: defina múltiplas assinaturas de
eval()para diferentes tipos de entradaAceita varargs (por exemplo,
eval(Integer...))Use tipos primitivos encapsulados (como
Integerem vez deint) para tratar entradasNULLcorretamente
Pré-requisitos
Antes de começar, verifique se você tem:
IntelliJ IDEA instalado
Apache Maven instalado
Acesso a um workspace do Realtime Compute for Apache Flink
Desenvolver uma UDSF
O Flink fornece exemplos de UDF com ambiente de desenvolvimento pré-configurado. Esses exemplos abrangem UDSFs, funções de agregação definidas pelo usuário (UDAFs) e funções de tabela definidas pelo usuário (UDTFs).
-
Baixe e descompacte o exemplo ASI_UDX_Demo em sua máquina local. Após a descompactação, a pasta
ASI_UDX-mainé criada com a seguinte estrutura:pom.xml: arquivo de configuração do projeto Maven. Define as coordenadas do projeto, dependências, regras de build e metadados relacionados.\ASI_UDX-main\src\main\java\ASI_UDF\ASI_UDF.java: implementação de exemplo da UDSF em Java.
O repositório ASI_UDX_Demo está hospedado em um site de terceiros. Podem ocorrer falhas ou atrasos no acesso.
No IntelliJ IDEA, clique em File > Open e selecione a pasta
ASI_UDX-main.-
Abra o arquivo
\ASI_UDX-main\src\main\java\ASI_UDF\ASI_UDF.javae atualize o métodoeval()com sua lógica personalizada. A implementação de exemplo extrai caracteres da posiçãobeginaté a posiçãoendde cada string de entrada:package ASI_UDF; import org.apache.flink.table.functions.ScalarFunction; public class ASI_UDF extends ScalarFunction { public String eval(String s, Integer begin, Integer end) { return s.substring(begin, end); } } -
Abra o arquivo
\ASI_UDX-main\pom.xmle configure as dependências do Maven correspondentes à sua versão do Flink. O exemplo a seguir mostra as principais dependências de pacotes JAR para o Flink 1.11. Se sua UDSF não depender de pacotes JAR adicionais, ignore esta etapa.<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>1.11.0</version> <!--<scope>provided</scope>--> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table</artifactId> <version>1.11.0</version> <type>pom</type> <!--<scope>provided</scope>--> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.11.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-common</artifactId> <version>1.11.0</version> </dependency> </dependencies> -
No diretório que contém o arquivo
pom.xml, execute o comando abaixo para empacotar o projeto:mvn package -Dcheckstyle.skipApós a conclusão do build, o arquivo
ASI_UDX-1.0-SNAPSHOT.jaré gerado no diretório\ASI_UDX-main\target\.
Registrar uma UDSF
Para registrar o pacote JAR como uma UDSF, consulte Gerencie funções definidas pelo usuário (UDFs).
Usar uma UDSF
Após o registro, chame a UDSF em um job Flink SQL.
-
Crie um job Flink SQL. Para orientações, consulte Mapa de desenvolvimento de jobs. O SQL de exemplo a seguir extrai os caracteres da segunda à quarta posição da string no campo
adeASI_UDSF_Sourcee grava os resultados emASI_UDSF_Sink: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; Na página Operation Center > Job O&M, localize seu job e clique em Start na coluna Actions. Após a inicialização do job, os caracteres da segunda à quarta posição do campo
aem cada linha deASI_UDSF_Sourcesão inseridos emASI_UDSF_Sink.