Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Python

Dernière mise à jour :Aug 09, 2026

Realtime Compute for Apache Flink prend en charge les fonctions définies par l'utilisateur (UDF) en Python pour les jobs Flink SQL. Créez des fonctions scalaires, d'agrégation et de table en Python, gérez les dépendances Python et optimisez les performances des UDF.

Types de fonctions définies par l'utilisateur

Catégorie

Description

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

Une UDSF associe zéro, une ou plusieurs valeurs scalaires à une nouvelle valeur scalaire. Elle traite une ligne d'entrée pour produire une valeur de sortie, établissant ainsi une correspondance un-à-un. Pour plus d'informations, consultez Fonctions scalaires définies par l'utilisateur (UDSF).

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

Une UDAF agrège plusieurs enregistrements en un seul, établissant une correspondance plusieurs-à-un. Pour plus d'informations, consultez Fonctions d'agrégation définies par l'utilisateur (UDAF).

fonction de table définie par l'utilisateur (UDTF)

Une UDTF accepte zéro, une ou plusieurs valeurs scalaires en tant que paramètres d'entrée. Contrairement à une fonction scalaire, une fonction de table peut renvoyer un nombre quelconque de lignes, chacune composée d'une ou plusieurs colonnes. Pour plus d'informations, consultez Fonctions de table définies par l'utilisateur (UDTF).

Utilisation des dépendances Python

Les clusters Realtime Compute for Apache Flink incluent des packages Python courants préinstallés, tels que Pandas, NumPy et PyArrow. Pour obtenir la liste des packages Python tiers préinstallés, consultez la rubrique Développement de jobs Python. Importez ces packages dans votre fonction avant utilisation, comme illustré dans l'exemple suivant.

@udf(result_type=DataTypes.FLOAT())
def percentile(values: List[float], percentile: float):
    import numpy as np
    return np.percentile(values, percentile)

Pour utiliser un package Python tiers non préinstallé, téléchargez-le en tant que fichier de dépendance lors de l'enregistrement de l'UDF Python. Pour plus d'informations, consultez les rubriques Gestion des fonctions définies par l'utilisateur (UDF) et Utilisation des dépendances Python.

Débogage du code

Utilisez le module de journalisation pour générer des journaux depuis votre fonction définie par l'utilisateur Python afin de faciliter le dépannage. L'exemple suivant illustre l'utilisation de la journalisation.

@udf(result_type=DataTypes.BIGINT())
def add(i, j):    
  logging.info("hello world")    
  return i + j

Une fois les journaux générés, consultez-les dans les fichiers journaux du TaskManager. Pour plus d'informations, voir la rubrique Affichage des journaux d'exécution.

Optimisation des performances

Préchargement des ressources

Préchargez les ressources lors de l'initialisation de la fonction pour éviter de les recharger à chaque appel de la méthode eval. Par exemple, chargez un modèle d'apprentissage profond volumineux une seule fois, puis effectuez des prédictions par lots.

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

class Predict(ScalarFunction):
    def open(self, function_context):
        import pickle

        with open("resources.zip/resources/model.pkl", "rb") as f:
            self.model = pickle.load(f)

    def eval(self, x):
        return self.model.predict(x)

predict = udf(Predict(), result_type=DataTypes.DOUBLE(), func_type="pandas")
Remarque

Pour savoir comment télécharger des fichiers de données Python, consultez la rubrique Utilisation des dépendances Python.

Fonctions définies par l'utilisateur asynchrones

Pour les scénarios intensifs en E/S, tels que l'accès à des bases de données externes ou les appels de services HTTP, utilisez des fonctions définies par l'utilisateur asynchrones. Une seule instance de fonction peut gérer plusieurs requêtes simultanément, répartissant le temps d'attente sur plusieurs appels et améliorant considérablement le débit du job. Cette fonctionnalité est prise en charge uniquement dans les versions VVR 11.7 et ultérieures et s'applique uniquement aux fonctions scalaires définies par l'utilisateur (UDSF). Pour plus de détails, consultez la rubrique Fonctions définies par l'utilisateur asynchrones.

Utilisation de la bibliothèque Pandas

Outre les fonctions définies par l'utilisateur Python standard, Realtime Compute for Apache Flink prend en charge les fonctions définies par l'utilisateur Pandas. Ces fonctions acceptent des structures de données Pandas telles que pandas.Series et pandas.DataFrame en entrée, ce qui permet d'exploiter des bibliothèques hautes performances comme Pandas et NumPy. Pour plus d'informations, consultez la documentation Vectorized User-defined Functions.

Paramètres

Les performances d'une fonction définie par l'utilisateur Python dépendent largement de son implémentation. En cas de problèmes de performances, commencez par optimiser la logique de la fonction. Les paramètres suivants influent également sur les performances.

Paramètre

Description

python.fn-execution.bundle.size

Les UDF Python s'exécutent de manière asynchrone. L'opérateur Java met les données en cache avant de les envoyer à un processus Python pour exécution. Lorsque le cache atteint un seuil, les données sont envoyées au processus Python. Le paramètre python.fn-execution.bundle.size contrôle le nombre maximal d'enregistrements pouvant être mis en cache.

La valeur par défaut est de 100 000 enregistrements.

python.fn-execution.bundle.time

Ce paramètre contrôle la durée maximale de mise en cache. Les données mises en cache sont envoyées pour traitement lorsque le nombre d'enregistrements atteint le seuil python.fn-execution.bundle.size ou lorsque la durée de mise en cache atteint le seuil python.fn-execution.bundle.time.

La valeur par défaut est de 1 000 millisecondes.

python.fn-execution.arrow.batch.size

Pour les UDF Pandas, ce paramètre spécifie le nombre maximal d'enregistrements qu'un lot Arrow peut contenir. La valeur par défaut est de 10 000.

Remarque

La valeur du paramètre python.fn-execution.arrow.batch.size ne peut pas être supérieure à la valeur du paramètre python.fn-execution.bundle.size.

Remarque

Définir des valeurs excessivement élevées pour ces paramètres peut s'avérer contre-productif. Une mise en tampon excessive pendant un point de contrôle peut entraîner une durée trop longue ou un échec des points de contrôle. Pour plus d'informations sur ces paramètres, consultez la documentation Configuration.

Rubriques connexes