Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Fonctions scalaires définies par l'utilisateur (UDSF) Python

Dernière mise à jour :Aug 13, 2026

Une fonction scalaire définie par l'utilisateur (UDSF) associe zéro, une ou plusieurs valeurs scalaires à une seule valeur scalaire. Chaque ligne en entrée génère exactement une valeur en sortie.

Cette rubrique explique comment créer, enregistrer et utiliser une UDSF Python dans Realtime Compute for Apache Flink.

Limites

Les contraintes suivantes s'appliquent au développement de fonctions définies par l'utilisateur (UDF) Python dans Realtime Compute for Apache Flink :

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 Prise en charge uniquement dans VVR 11.9-preview1 et versions ultérieures.
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é.

Méthodes de développement

Vous disposez de deux méthodes pour développer une UDSF Python :

  • Empaquetez votre code Python, téléchargez le package, puis enregistrez-le en tant que fonction sur la plateforme.

  • Déclarez la logique du code sous forme de fonction inline dans une instruction SQL.

Si la logique de votre UDSF est simple, privilégiez le développement sous forme de fonction inline.

Créer une UDSF

Les étapes suivantes prennent Windows comme environnement d'exemple. Flink fournit un référentiel d'exemples qui inclut des implémentations pour les UDSF, les fonctions d'agrégation définies par l'utilisateur (UDAF) et les fonctions table définies par l'utilisateur (UDTF).
  1. Téléchargez et décompressez python_demo-master sur votre machine locale.

    Il s'agit d'un référentiel GitHub tiers. L'accès peut être lent ou intermittent.
  2. Dans PyCharm, choisissez File > Open et ouvrez le répertoire python_demo-master décompressé.

  3. Ouvrez udfs.py dans le chemin \python_demo-master\udx et définissez votre UDSF.

    from pyflink.table import DataTypes
    from pyflink.table.udf import udf
    
    @udf(result_type=DataTypes.STRING())
    def sub_string(s: str, begin: int, end: int):
        return s[begin:end]

    L'exemple sub_string extrait les caractères de la position begin à la position end dans la chaîne d'entrée.

  4. Depuis le répertoire \python_demo-master, exécutez la commande suivante pour empaqueter le répertoire udx :

    zip -r python_demo.zip udx

    Lorsque python_demo.zip apparaît dans \python_demo-master\, le package est prêt.

Enregistrer une UDSF

Une fois le package créé, enregistrez l'UDSF dans la console Realtime Compute for Apache Flink. Pour connaître les étapes d'enregistrement, consultez Gérer les fonctions définies par l'utilisateur (UDF).

Utiliser une UDSF

Après avoir enregistré l'UDSF, utilisez-la dans un job Flink SQL.

  1. Créez un brouillon à l'aide de Flink SQL. Pour plus de détails, consultez Présentation du développement de jobs. L'exemple suivant appelle ASI_UDSF (le nom enregistré de votre UDSF) pour extraire les caractères des positions 2 à 4 du champ a dans la table source :

    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;
  2. Dans le volet de navigation de gauche de la console de développement, choisissez O&M > Deployments. Localisez le déploiement, puis cliquez sur Start dans la colonne Actions. Une fois le déploiement démarré, les caractères aux positions 2 à 4 du champ a dans ASI_UDSF_Source sont écrits dans ASI_UDSF_Sink.

Fonctions scalaires inline Python

Une fonction inline intègre son implémentation directement dans une instruction CREATE FUNCTION, ce qui permet de définir et d'enregistrer la fonction dans la même instruction SQL. L'exemple ci-dessous définit une fonction inline qui masque les adresses e-mail. Déclarez l'intégralité de la logique du code entre les délimiteurs $$.

CREATE TEMPORARY FUNCTION mask_email(email STRING)
RETURNS STRING AS $$
if email is None:
    return None

name, separator, domain = email.partition("@")
if not separator:
    return email

return name[:1] + "***@" + domain
$$ LANGUAGE PYTHON;
Remarque
  1. Python est sensible à l'indentation. Commencez le code de premier niveau du corps de la fonction à la colonne 0 et maintenez une indentation relative correcte au sein des blocs de code.

  2. Les fonctions inline ne prennent pas en charge l'exécution vectorisée et ne peuvent pas être déclarées comme fonctions non déterministes.

  3. Seules les fonctions TEMPORARY sont prises en charge. Les définitions de fonctions ne peuvent pas encore être persistées sur la plateforme.

Fonctions définies par l'utilisateur asynchrones

Pour les UDF effectuant des opérations intensives en E/S, telles que l'accès à des bases de données externes ou des requêtes HTTP, utilisez des fonctions définies par l'utilisateur asynchrones. Une seule UDF asynchrone peut traiter plusieurs requêtes d'E/S simultanément, répartissant ainsi le temps d'attente entre les requêtes et améliorant le débit du job.

Limites

  • Prise en charge uniquement sur VVR 11.7 et versions ultérieures. VVR PyFlink 11.7 ou version ultérieure est requis. Pour plus de détails, consultez ververica-flink.

    pip3 install "ververica-flink>=11.7"

  • Seules les fonctions scalaires définies par l'utilisateur (UDSF) asynchrones sont prises en charge.

  • Seul le mode processus Python est pris en charge, c'est-à-dire python.execution-mode=process.

  • Les fonctions définies par l'utilisateur asynchrones Pandas ne sont pas encore prises en charge.

  • Les fonctions inline ne sont pas encore prises en charge.

Utilisation

Une fonction définie par l'utilisateur asynchrone peut être implémentée sous forme de fonction async Python ou comme sous-classe de la classe de fonction asynchrone. Le code d'exemple est présenté ci-dessous.

import asyncio

from pyflink.table import DataTypes
from pyflink.table.udf import AsyncScalarFunction, udf

# Method 1: Use a Python async function
@udf(result_type=DataTypes.STRING())
async def async_api_call(product_id: str) -> str:
    await asyncio.sleep(0.05)
    return f"product_{product_id}"

# Method 2: Subclass the asynchronous function class
class AsyncUserLookup(AsyncScalarFunction):
    def open(self, function_context):
        self.cache = {}

    async def eval(self, user_id: str) -> str:
        if user_id in self.cache:
            return self.cache[user_id]

        await asyncio.sleep(0.05)
        result = f"user_{user_id}"
        self.cache[user_id] = result
        return result

    def close(self):
        self.cache.clear()

async_user_lookup = udf(
    AsyncUserLookup(),
    input_types=[DataTypes.STRING()],
    result_type=DataTypes.STRING()
)

Les fonctions définies par l'utilisateur asynchrones s'enregistrent et s'utilisent de la même manière que les fonctions synchrones.

Paramètres de configuration

Les paramètres suivants contrôlent le comportement d'exécution des fonctions définies par l'utilisateur asynchrones.

Paramètre

Valeur par défaut

Description

table.exec.async-scalar.max-concurrent-operations

10

Nombre maximal d'appels asynchrones simultanés par instance d'opérateur. Valeur par défaut : 10.

table.exec.async-scalar.timeout

3 min

Délai d'expiration pour un appel asynchrone unique.

table.exec.async-scalar.retry-strategy

FIXED_DELAY

Stratégie de nouvelle tentative après l'échec d'un appel asynchrone. Valeurs prises en charge :

  • FIXED_DELAY : Nouvelle tentative après un temps d'attente fixe.

  • NO_RETRY : Aucune nouvelle tentative.

table.exec.async-scalar.retry-delay

100 ms

Temps d'attente pour les nouvelles tentatives à délai fixe.

Remarque

Effectif uniquement lorsque table.exec.async-scalar.retry-strategy est défini sur FIXED_DELAY.

table.exec.async-scalar.max-attempts

3

Nombre maximal de tentatives avant qu'un appel asynchrone ne soit considéré comme échoué.

Remarque

Effectif uniquement lorsque table.exec.async-scalar.retry-strategy est défini sur FIXED_DELAY.