Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Java

Dernière mise à jour :Aug 09, 2026

Realtime Compute for Apache Flink prend en charge les fonctions définies par l'utilisateur (UDF) Java dans les jobs Flink SQL. Découvrez les types d'UDF, le passage de paramètres et les notes d'utilisation importantes.

Notes importantes

  • Pour éviter les conflits de dépendances de packages JAR lors du développement des UDF, tenez compte des points suivants :

    • Assurez-vous que la version de Flink sélectionnée sur la page de développement SQL correspond à la version de Flink indiquée dans la dépendance Pom.

    • Pour les dépendances liées à Flink, définissez l'étendue sur provided en ajoutant <scope>provided</scope>.

    • Emblez les autres dépendances tierces à l'aide de la méthode Shade. Pour plus d'informations, consultez le Plug-in Apache Maven Shade.

    Pour plus d'informations sur les conflits de dépendances Flink, consultez la rubrique Comment résoudre les conflits de dépendances Flink ?.

  • Pour prévenir les délais d'expiration causés par des appels fréquents à une UDF dans un job SQL, vous pouvez télécharger le package JAR de l'UDF en tant que fichier de dépendance. Déclarez ensuite la fonction dans le job à l'aide de la syntaxe CREATE TEMPORARY FUNCTION. Par exemple : CREATE TEMPORARY FUNCTION 'GetJson' AS 'com.soybean.function.udf.MyGetJsonObjectFunction';

Classifications des UDF

Flink prend en charge les trois types d'UDF suivants.

Classification

Description

Fonction scalaire définie par l'utilisateur (UDSF)

Une UDSF mappe zéro, une ou plusieurs valeurs scalaires vers une nouvelle valeur scalaire. Elle établit une relation un-à-un entre l'entrée et la sortie. Cela signifie qu'elle lit une ligne de données et écrit une seule valeur de sortie. Pour plus d'informations, consultez la rubrique Fonctions scalaires définies par l'utilisateur (UDSF).

Fonction d'agrégation définie par l'utilisateur (UDAF)

Une UDAF agrège plusieurs enregistrements d'entrée en une seule valeur de sortie. Pour plus d'informations, consultez la rubrique Fonctions d'agrégation définies par l'utilisateur (UDAF).

Fonction table définie par l'utilisateur (UDTF)

Une UDTF accepte zéro, une ou plusieurs valeurs scalaires comme paramètres d'entrée. Les paramètres peuvent être de longueur variable. Elle est similaire à une UDSF, mais peut renvoyer un nombre quelconque de lignes en sortie, et pas seulement une seule valeur. Les lignes renvoyées peuvent comporter une ou plusieurs colonnes. Un seul appel de fonction peut produire plusieurs lignes ou colonnes. Pour plus d'informations, consultez la rubrique Fonctions table définies par l'utilisateur (UDTF).

Enregistrement des UDF

  • Pour savoir comment enregistrer une UDF globale, consultez la rubrique UDF globales.

  • Pour savoir comment enregistrer une UDF au niveau du job, consultez la rubrique UDF au niveau du job.

Passage de paramètres à une UDF

Vous pouvez configurer des paramètres pour une UDF dans la console de développement Flink et y faire référence dans le code de l'UDF, ce qui vous permet de modifier rapidement les valeurs des paramètres directement dans la console.

Les UDF proposent une méthode optionnelle open(FunctionContext context). FunctionContext transmet les éléments de configuration personnalisés sous forme de paramètres. Le processus est le suivant :

  1. Dans l'onglet Configuration de la page O&M > Deployments de la console de développement Flink, ajoutez l'élément de configuration pipeline.global-job-parameters dans la section Other Configuration de Parameters.

    pipeline.global-job-parameters: | 
      'k1:{hi,hello}',
      'k2:"str:ing,str:ing"',
      'k3:"str""ing,str:ing"'

    FunctionContext#getJobParameter ne peut récupérer que les valeurs définies dans pipeline.global-job-parameters. Vous devez inclure tous les éléments de configuration utilisés par l'UDF dans ce paramètre. Les étapes suivantes décrivent comment configurer cet élément.

    Étape

    Action

    Procédure

    Exemple

    Étape 1

    Définissez des paires clé-valeur.

    Séparez la clé et la valeur par deux-points (:) et placez chaque paire clé-valeur entre guillemets simples (').

    Remarque
    • Si une clé ou une valeur contient un deux-points (:), placez-la entre guillemets doubles (").

    • Si une clé ou une valeur contient un deux-points (:) ou un guillemet double ("), vous devez l'échapper avec deux guillemets doubles consécutifs ("").

    • Si key = k1 et value = {hi,hello}, définissez la paire comme 'k1:{hi,hello}'.

    • Si key = k2 et value = str:ing,str:ing, définissez la paire comme 'k2:"str:ing,str:ing"'

    • Si key = k3 et value = str"ing,str:ing, définissez la paire comme 'k3:"str""ing,str:ing"'

    Étape 2

    Formatez la valeur finale de pipeline.global-job-parameters sous forme de fichier YAML.

    Placez chaque paire clé-valeur sur une nouvelle ligne et séparez-les par des virgules (,).

    Remarque
    • Une chaîne multiligne dans un fichier YAML commence par une barre verticale (|).

    • Chaque ligne d'une chaîne multiligne dans un fichier YAML doit avoir la même indentation.

    pipeline.global-job-parameters: | 
      'k1:{hi,hello}',
      'k2:"str:ing,str:ing"',
      'k3:"str""ing,str:ing"'
  2. Dans le code de l'UDF, utilisez FunctionContext#getJobParameter pour récupérer la valeur de chaque élément. L'extrait de code suivant illustre cette approche.

    Exemple :

    context.getJobParameter("k1", null); // Returns the string {hi,hello}.
    context.getJobParameter("k2", null); // Returns the string str:ing,str:ing.
    context.getJobParameter("k3", null); // Returns the string str"ing,str:ing.
    context.getJobParameter("pipeline.global-job-parameters", null); // Returns null. You can only get the content defined in pipeline.global-job-parameters, not any other job configuration item.

Paramètres nommés

Remarque

Seuls Ververica Runtime (VVR) 8.0.7 et versions ultérieures prennent en charge l'utilisation de paramètres nommés pour implémenter des UDF.

Les appels de fonction standard exigent que tous les paramètres soient fournis dans le bon ordre, ce qui est source d'erreurs pour les fonctions comportant de nombreux paramètres et n'autorise pas l'omission des paramètres facultatifs. Les paramètres nommés vous permettent de spécifier uniquement les paramètres dont vous avez besoin. L'exemple ScalarFunction suivant montre comment utiliser les paramètres nommés.

// Implement a user-defined scalar function. The last two input parameters are optional (isOptional = true).
public class MyFuncWithNamedArgs extends ScalarFunction  {
	private static final long serialVersionUID = 1L;

	public String eval(@ArgumentHint(name = "f1", isOptional = false, type = @DataTypeHint("STRING")) String f1,
			@ArgumentHint(name = "f2", isOptional = true, type = @DataTypeHint("INT")) Integer i2,
			@ArgumentHint(name = "f3", isOptional = true, type = @DataTypeHint("LONG")) Long l3) {

		if (i2 != null) {
			return "i2#" + i2;
		}
		if (l3 != null) {
			return "l3#" + l3;
		}
		return "default#" + f1;
	}
}

Lorsque vous utilisez cette UDF en SQL, vous pouvez spécifier uniquement le paramètre obligatoire ou inclure sélectivement les paramètres facultatifs.

CREATE TEMPORARY FUNCTION MyNamedUdf AS 'com.aliyun.example.MyFuncWithNamedArgs';

CREATE temporary TABLE s1 (
    a INT,
    b BIGINT,
    c VARCHAR,
    d VARCHAR,
    PRIMARY KEY(a) NOT ENFORCED
) WITH (
    'connector' = 'datagen',
    'rows-per-second'='1'
);

CREATE temporary TABLE sink (
    a INT,
    b VARCHAR,
    c VARCHAR,
    d VARCHAR
) WITH (
    'connector' = 'print'
);

INSERT INTO sink
SELECT a,
    -- Specify only the first required parameter
    MyNamedUdf(f1 => c) arg1_res,
    -- Specify the first required parameter and the second optional parameter
    MyNamedUdf(f1 => c, f2 => a) arg2_res,
    -- Specify the first required parameter and the third optional parameter
    MyNamedUdf(f1 => c, f3 => d) arg3_res
FROM s1;

Références