Une fonction scalaire définie par l'utilisateur (UDSF) transforme zéro, une ou plusieurs valeurs scalaires en une nouvelle valeur scalaire, établissant une relation un-à-un entre les lignes d'entrée et les valeurs de sortie. Utilisez une UDSF lorsque les fonctions Flink SQL intégrées ne permettent pas d'exprimer votre logique personnalisée.
Pour obtenir un aperçu de tous les types de fonctions définies par l'utilisateur (UDF), consultez la page Fonctions définies par l'utilisateur .
Fonctionnement
Une UDSF étend la classe ScalarFunction d'Apache Flink et implémente une ou plusieurs méthodes eval(). Flink appelle la méthode eval() une fois pour chaque ligne d'entrée et utilise la valeur renvoyée comme scalaire de sortie.
public class SubstringFunction extends ScalarFunction {
public String eval(String s, Integer begin, Integer end) {
return s.substring(begin, end);
}
}
Exigences relatives à la méthode eval() :
Déclarez la méthode comme
public.La surcharge de méthode est prise en charge : définissez plusieurs signatures
eval()pour différents types d'entrée.Les arguments variables sont pris en charge (par exemple,
eval(Integer...)).Utilisez des types primitifs encapsulés (par exemple,
Integerau lieu deint) afin de gérer correctement les entréesNULL.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
IntelliJ IDEA installé
Apache Maven installé
Accès à un espace de travail Realtime Compute for Apache Flink
Développer une UDSF
Flink fournit des exemples d'UDF qui préconfigurent l'environnement de développement. Ces exemples couvrent les UDSF, les fonctions d'agrégation définies par l'utilisateur (UDAF) et les fonctions tabulaires définies par l'utilisateur (UDTF).
-
Téléchargez et décompressez l'exemple ASI_UDX_Demo sur votre machine locale. Après la décompression, le dossier
ASI_UDX-mainest créé avec la structure suivante :pom.xml: fichier de configuration du projet Maven. Il définit les coordonnées du projet, les dépendances, les règles de build et les métadonnées associées.\ASI_UDX-main\src\main\java\ASI_UDF\ASI_UDF.java: exemple d'implémentation Java de la UDSF.
L'exemple ASI_UDX_Demo est hébergé sur un site tiers. Vous pouvez rencontrer des échecs d'accès ou des délais.
Dans IntelliJ IDEA, cliquez sur File > Open et sélectionnez le dossier
ASI_UDX-main.-
Ouvrez le fichier
\ASI_UDX-main\src\main\java\ASI_UDF\ASI_UDF.javaet mettez à jour la méthodeeval()avec votre logique personnalisée. L'implémentation d'exemple extrait les caractères situés entre la positionbeginet la positionendde chaque chaîne d'entrée :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); } } -
Ouvrez le fichier
\ASI_UDX-main\pom.xmlet configurez les dépendances Maven pour votre version de Flink. L'exemple suivant présente les principales dépendances de package JAR pour Flink 1.11. Si votre UDSF ne dépend pas de packages JAR supplémentaires, ignorez cette étape.<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> -
Dans le répertoire contenant le fichier
pom.xml, exécutez la commande suivante pour empaqueter le projet :mvn package -Dcheckstyle.skipUne fois la build terminée, le fichier
ASI_UDX-1.0-SNAPSHOT.jarest créé dans le répertoire\ASI_UDX-main\target\.
Enregistrer une UDSF
Pour enregistrer le package JAR en tant que UDSF, consultez la rubrique Gérer les fonctions définies par l'utilisateur (UDF).
Utiliser une UDSF
Après l'enregistrement, appelez la UDSF dans une tâche Flink SQL.
-
Créez une tâche Flink SQL. Pour obtenir des instructions, consultez la page Carte de développement des tâches. L'exemple SQL suivant extrait les caractères situés entre la deuxième et la quatrième position de la chaîne dans le champ
adeASI_UDSF_Sourceet écrit les résultats dansASI_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; Sur la page Operation Center > Job O&M, recherchez votre tâche et cliquez sur Start dans la colonne Actions. Une fois la tâche démarrée, les caractères situés entre la deuxième et la quatrième position du champ
ade chaque ligne deASI_UDSF_Sourcesont insérés dansASI_UDSF_Sink.