Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:UDAF Python

Dernière mise à jour :Aug 09, 2026

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 :

  1. create_accumulator() — crée un accumulateur vide pour contenir l'état initial.

  2. accumulate(accumulator, ...) — appelé une fois par ligne d'entrée pour mettre à jour l'accumulateur.

  3. 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.
  1. Téléchargez et décompressez le package python_demo-master sur votre machine.

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

  3. Ouvrez le fichier \python_demo-master\udx\udfs.py et modifiez-le pour qu'il corresponde à votre logique métier. L'exemple ci-dessous définit weighted_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())]))
  4. Depuis le répertoire \python_demo-master\ , archivez le dossier udx :

    zip -r python_demo.zip udx

    La UDAF est créée lorsque le fichier python_demo.zip apparaî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

  1. 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 a dans ASI_UDAF_Source, en utilisant le champ b comme 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;
  2. 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 champ b comme poids — est écrite dans chaque ligne de ASI_UDAF_Sink .