Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:UDTF Python

Dernière mise à jour :Aug 19, 2026

Cette rubrique explique comment créer, enregistrer et utiliser une fonction table-valued définie par l'utilisateur (UDTF) Python dans Realtime Compute for Apache Flink.

Description

Une UDTF accepte zéro, un ou plusieurs paramètres scalaires en entrée. Ces paramètres peuvent être de longueur variable. Contrairement aux fonctions classiques qui retournent une valeur unique, les UDTF peuvent renvoyer un nombre quelconque de lignes comportant une ou plusieurs colonnes. Chaque appel à une UDTF produit ainsi plusieurs lignes ou colonnes. Bien que similaires aux fonctions scalaires définies par l'utilisateur (UDF), les UDTF s'en distinguent par la structure de leurs résultats.

Limites

Le développement de fonctions définies par l'utilisateur (UDF) Python dans Realtime Compute for Apache Flink est soumis aux contraintes suivantes :

Contrainte Exigence
Version d'Apache Flink 1.12 et versions ultérieures
Version de Python Préinstallée sur chaque espace de travail. VVR antérieure à 8.0.11 : Python 3.7.9. VVR 8.0.11 et versions ultérieures : Python 3.9.21.
Version du JDK JDK 8 et JDK 11. Les packages JAR tiers doivent être compatibles avec JDK 8 ou JDK 11.
Version de Scala Uniquement Scala 2.11 open source. Les packages JAR tiers doivent être compatibles avec Scala 2.11.
Fonctions inline Prises en charge uniquement à partir de VVR 11.9-preview1.
Important

Après une mise à niveau vers VVR 8.0.11 ou version ultérieure, testez, déployez et exécutez à nouveau vos brouillons PyFlink existants afin de confirmer leur compatibilité.

Créer une UDTF

Remarque

Flink fournit des exemples de code d'extensions définies par l'utilisateur (UDX) Python pour faciliter le développement. Ces exemples couvrent l'implémentation de UDF Python, de fonctions d'agrégation définies par l'utilisateur (UDAF) Python et de UDTF Python. Cette section détaille la création d'une UDTF sous le système d'exploitation Windows.

Important

Le nom enregistré d'une fonction personnalisée ne peut contenir que des lettres minuscules, des chiffres et des traits d'union (-). Les caractères de soulignement (_) et autres caractères spéciaux ne sont pas pris en charge.

Le nom de la fonction ne doit pas entrer en conflit avec celui d'une fonction intégrée. L'environnement Flink géré comporte déjà certaines fonctions intégrées préenregistrées, telles que split. Si une fonction personnalisée porte le même nom qu'une fonction intégrée, l'enregistrement échoue ou une erreur de validation SQL survient. Il est recommandé d'utiliser un préfixe personnalisé pour les noms de fonctions, par exemple my-split.

  1. Téléchargez et décompressez python_demo-master sur votre machine locale.

  2. Dans PyCharm, choisissez File > Open et ouvrez le répertoire python_demo-master décompressé.

  3. Double-cliquez sur le fichier udtfs.py situé dans le répertoire \python_demo-master\udx. Modifiez ensuite le contenu du fichier selon vos besoins métier.

    Dans cet exemple, my_split définit le code permettant de séparer une chaîne de caractères en plusieurs colonnes à l'aide de barres verticales (|).

    from pyflink.table import DataTypes
    from pyflink.table.udf import udtf
    
    @udtf(result_types=[DataTypes.STRING(), DataTypes.STRING()])
    def my_split(s: str):
        splits = s.split("|")
        yield splits[0], splits[1]
  4. Accédez au répertoire \python_demo contenant le dossier udx, puis exécutez la commande suivante pour empaqueter les fichiers du répertoire :

    zip -r python_demo.zip udx

    La présence du package python_demo.zip dans le répertoire \python_demo\ confirme que le développement de la UDTF est terminé.

Enregistrer une UDTF

Pour plus d'informations sur l'enregistrement d'une UDTF, consultez Gérer les UDF.

Utiliser une UDTF

Une fois la UDTF enregistrée, suivez les étapes ci-dessous pour l'utiliser :

  1. Créez un brouillon à l'aide de Flink SQL. Pour plus d'informations, consultez Présentation du développement de jobs.

    L'exemple suivant illustre la concaténation de la chaîne « aa » avec le champ message de chaque ligne de la table ASI_UDTF_Source via une barre verticale (|). La chaîne résultante est ensuite scindée en plusieurs colonnes par ce même séparateur :

    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(my_split(concat_ws('|', `message`, 'aa'))) as T(name,place);
  2. Dans le volet de navigation de gauche de la console de développement Realtime Compute for Apache Flink, choisissez O&M > Deployments. Dans la page Deployments, localisez le déploiement souhaité et cliquez sur Start dans la colonne Actions.

    Une fois le déploiement démarré, deux colonnes de données sont insérées dans la table ASI_UDTF_Sink. Celles-ci contiennent les chaînes concaténées séparées par des barres verticales (|).