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")
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 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 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 |
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
Pour savoir comment enregistrer, mettre à jour et supprimer une fonction définie par l'utilisateur, consultez la rubrique Gestion des fonctions définies par l'utilisateur (UDF).
Pour consulter des démonstrations sur le développement et l'utilisation de fonctions définies par l'utilisateur Python, reportez-vous aux rubriques Fonctions d'agrégation définies par l'utilisateur (UDAF), Fonctions scalaires définies par l'utilisateur (UDSF) et Fonctions de table définies par l'utilisateur (UDTF).
Pour apprendre à utiliser des environnements virtuels Python personnalisés, des packages Python tiers, des packages JAR et des fichiers de données dans un job Flink Python, consultez la rubrique Utilisation des dépendances Python.
Pour consulter des démonstrations sur le développement et l'utilisation de fonctions définies par l'utilisateur Java, reportez-vous aux rubriques Fonctions d'agrégation définies par l'utilisateur (UDAF), Fonctions scalaires définies par l'utilisateur (UDSF) et Fonctions de table définies par l'utilisateur (UDTF).
Pour savoir comment déboguer et optimiser les fonctions définies par l'utilisateur Java, consultez la rubrique Présentation.
Implémentation du tri et de l'agrégation des données à l'aide d'une UDAF.