Une fonction d'agrégation définie par l'utilisateur (UDAF) agrège plusieurs valeurs d'entrée en une seule valeur de sortie, en associant plusieurs lignes d'entrée à un résultat unique par groupe.
Cette rubrique explique comment créer, enregistrer et utiliser une UDAF Python dans Realtime Compute for Apache Flink.
Limites
Apache Flink 1.12 ou version ultérieure est requis.
-
Python est préinstallé sur l'espace de travail Realtime Compute for Apache Flink. Écrivez votre code avec la version de Python préinstallée.
Python 3.7.9 est préinstallé pour les versions de Ververica Runtime (VVR) antérieures à la 8.0.11. Python 3.9.21 est préinstallé pour VVR 8.0.11 ou version ultérieure. Après la mise à niveau vers VVR 8.0.11 ou version ultérieure, retestez, redéployez et réexécutez tous les jobs PyFlink construits sur une version VVR antérieure.
JDK 8 et JDK 11 sont pris en charge dans l'environnement d'exécution. Si le déploiement de votre application Python dépend d'un fichier JAR tiers, assurez-vous que ce fichier JAR est compatible avec JDK 8 ou JDK 11.
Seule la version open source Scala 2.11 est prise en charge. Si le déploiement de votre application Python dépend d'un fichier JAR tiers, assurez-vous que ce fichier JAR est compatible avec Scala 2.11.
Fonctionnement
Une UDAF utilise un accumulateur pour suivre l'état intermédiaire de l'agrégation entre les lignes d'entrée. L'accumulateur est créé une fois par groupe et mis à jour ligne par ligne jusqu'au traitement de toutes les entrées.
L'ordre d'exécution pour chaque groupe est le suivant :
create_accumulator()— crée un accumulateur vide pour contenir l'état initial.accumulate(accumulator, ...)— appelé une fois par ligne d'entrée pour mettre à jour l'accumulateur.get_value(accumulator)— appelé après le traitement de toutes les lignes du groupe pour renvoyer le résultat final.
Référence des méthodes
| Méthode | Requise | Quand l'implémenter |
|---|---|---|
create_accumulator() |
Oui | Toujours |
accumulate(...) |
Oui | Toujours |
get_value(...) |
Oui | Toujours |
retract(...) |
Conditionnelle | Lorsque l'agrégation peut recevoir des messages de rétraction. |
Créer une UDAF
Flink fournit des exemples de code pour les extensions définies par l'utilisateur (UDX), y compris les UDAF, les fonctions définies par l'utilisateur (UDF) et les fonctions de table définies par l'utilisateur (UDTF). Les étapes ci-dessous utilisent un environnement Windows.
Téléchargez et décompressez le package python_demo-master sur votre machine.
Dans PyCharm, choisissez File > Open et ouvrez le répertoire
python_demo-masterdécompressé.-
Ouvrez le fichier
\python_demo-master\udx\udfs.pyet modifiez-le pour qu'il corresponde à votre logique métier. L'exemple ci-dessous définitweighted_avg, qui calcule une moyenne pondérée sur les données actuelles et historiques.from pyflink.common import Row from pyflink.table import AggregateFunction, DataTypes from pyflink.table.udf import udaf class WeightedAvg(AggregateFunction): def create_accumulator(self): # Row(sum, count) return Row(0, 0) def get_value(self, accumulator: Row) -> float: if accumulator[1] == 0: return 0 else: return accumulator[0] / accumulator[1] def accumulate(self, accumulator: Row, value, weight): accumulator[0] += value * weight accumulator[1] += weight def retract(self, accumulator: Row, value, weight): accumulator[0] -= value * weight accumulator[1] -= weight weighted_avg = udaf(f=WeightedAvg(), result_type=DataTypes.DOUBLE(), accumulator_type=DataTypes.ROW([ DataTypes.FIELD("f0", DataTypes.BIGINT()), DataTypes.FIELD("f1", DataTypes.BIGINT())])) -
Depuis le répertoire
\python_demo-master\, archivez le dossierudx:zip -r python_demo.zip udxLa UDAF est créée lorsque le fichier
python_demo.zipapparaît dans le répertoire\python_demo-master\.
Enregistrer une UDAF
Pour les étapes d'enregistrement, consultez Gérer les fonctions définies par l'utilisateur (UDF).
Utiliser une UDAF
-
Créez un brouillon Flink SQL. Pour plus de détails, consultez Présentation du développement de jobs. L'exemple suivant calcule la moyenne pondérée du champ
adansASI_UDAF_Source, en utilisant le champbcomme poids.CREATE TEMPORARY TABLE ASI_UDAF_Source ( a BIGINT, b BIGINT ) WITH ( 'connector' = 'datagen' ); CREATE TEMPORARY TABLE ASI_UDAF_Sink ( avg_value DOUBLE ) WITH ( 'connector' = 'blackhole' ); INSERT INTO ASI_UDAF_Sink SELECT weighted_avg(a, b) FROM ASI_UDAF_Source; Dans le volet de navigation de gauche de la console de développement Realtime Compute for Apache Flink, choisissez O&M > Deployments. Recherchez le déploiement cible et cliquez sur Start dans la colonne Actions. Une fois le déploiement démarré, la moyenne pondérée du champ
a— avec le champbcomme poids — est écrite dans chaque ligne deASI_UDAF_Sink.