Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:UDAF Java

Dernière mise à jour :Aug 20, 2026

Une fonction d'agrégation définie par l'utilisateur (UDAF) réduit plusieurs lignes d'entrée en une seule valeur de sortie, établissant ainsi une correspondance plusieurs-vers-un. Utilisez les UDAF lorsque les fonctions d'agrégation intégrées telles que SUM ou MAX ne couvrent pas votre logique d'agrégation.

Pour obtenir des informations générales sur les fonctions définies par l'utilisateur dans Flink, consultez la documentation Fonctions définies par l'utilisateur.

Les ressources Fonctions définies par l'utilisateur et ASI_UDX_Demo sont hébergées sur des sites tiers et peuvent parfois être lentes ou indisponibles.

Fonctionnement

Une UDAF utilise un accumulateur pour stocker l'état intermédiaire de l'agrégation. Pour chaque groupe de lignes partageant la même clé GROUP BY, le runtime appelle trois méthodes séquentiellement :

  1. createAccumulator() : crée un nouvel accumulateur avec un état initial.

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

  3. getValue(acc) : appelée après le traitement de toutes les lignes pour renvoyer le résultat final.

Par exemple, si une table contient une colonne numérique et que vous souhaitez calculer la somme cumulative des valeurs 1, 2 et 3 au sein du même groupe :

  • Le runtime appelle createAccumulator() une fois pour initialiser l'accumulateur avec sum = 0.

  • Il appelle accumulate() pour chaque ligne, mettant à jour sum à 1, puis à 3, et enfin à 6.

  • Il appelle getValue() pour renvoyer le résultat final.

La sortie dépend de l'activation ou non du mini-batch :

  • Sans mini-batch (par défaut) : une sortie par ligne, soit 1, 3, 6.

  • Avec mini-batch activé : seul le résultat final est émis, soit 6. Le nombre de sorties intermédiaires varie selon les paramètres du mini-batch et la distribution des données.

Pour la configuration du mini-batch, consultez la rubrique Optimisation de Flink SQL.

Méthodes requises

Toute implémentation de AggregateFunction doit définir les trois méthodes suivantes :

Méthode Objectif
createAccumulator() Renvoie un nouvel accumulateur avec un état initial
accumulate(acc, ...) Met à jour l'accumulateur avec une ligne d'entrée
getValue(acc) Renvoie le résultat agrégé final

Des méthodes supplémentaires sont disponibles selon votre cas d'utilisation :

Méthode Objectif
retract(acc, ...) Prend en charge la rétraction d'un message généré par un opérateur en amont
merge(acc, iterable) Prend en charge l'optimisation d'agrégation en deux étapes local-global

Création d'une UDAF

Realtime Compute for Apache Flink fournit un projet de démonstration UDF (ASI_UDX_Demo) avec un environnement de développement préconfiguré, ce qui évite toute configuration préalable de l'environnement.

Le projet de démonstration inclut des exemples d'implémentation pour les fonctions scalaires définies par l'utilisateur (UDSF), les UDAF et les fonctions tabulaires définies par l'utilisateur (UDTF).

Prérequis

Avant de commencer, assurez-vous de disposer des éléments suivants :

  • IntelliJ IDEA installé

  • Maven installé

  • Environnement de développement Java configuré

Étapes

  1. Téléchargez et décompressez le projet ASI_UDX_Demo sur votre machine locale. Le dossier décompressé ASI_UDX-main contient :

    • pom.xml : configuration du projet Maven, incluant les coordonnées, les dépendances et les règles de build.

    • \ASI_UDX-main\src\main\java\ASI_UDAF\ASI_UDAF.java : exemple d'implémentation d'UDAF.

  2. Dans IntelliJ IDEA, cliquez sur File > Open, puis sélectionnez le dossier ASI_UDX-main.

  3. Ouvrez le fichier pom.xml situé dans le répertoire \ASI_UDX-main\ et configurez les dépendances. Le fichier inclut la dépendance minimale pour Flink 1.12 :

    • Si votre job n'a aucune dépendance supplémentaire, passez à l'étape suivante.

    • Si votre job nécessite des dépendances supplémentaires, ajoutez-les au fichier pom.xml.

    <dependencies>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-common</artifactId>
            <version>1.12.7</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>

    Utilisez la dernière version mineure de la version majeure d'Apache Flink correspondant à votre version Ververica Runtime (VVR). Pour les correspondances entre les versions VVR et Flink, consultez la rubrique Vue d'ensemble.

  4. Ouvrez le fichier \ASI_UDX-main\src\main\java\ASI_UDAF\ASI_UDAF.java et implémentez votre logique d'agrégation. L'exemple implémente une sommation cumulative :

    package ASI_UDAF;
    
    import org.apache.flink.table.functions.AggregateFunction;
    
    import java.util.Iterator;
    
    public class ASI_UDAF{
        public static class AccSum{
            public long sum;
        }
    
        public static class MySum extends AggregateFunction<Long, AccSum>{
    
            @Override
            public Long getValue(AccSum acSum){
                return acSum.sum;
            }
    
            @Override
            public AccSum createAccumulator(){
                AccSum acCount= new AccSum();
                acCount.sum=0;
                return acCount;
            }
    
            public void accumulate(AccSum acc,long num){
                acc.sum += num;
            }
    
            /**
            * Supports retracting a message generated by an upstream operator.
            */
            public void retract(AccSum acc,long num){
                acc.sum -= num;
            }
    
            /**
            * Supports local-global two-stage aggregate optimization.
            */
            public void merge(AccSum acc,Iterable<AccSum> it){
                Iterator<AccSum> iter=it.iterator();
                while(iter.hasNext()){
                    AccSum accSum=iter.next();
                    if(null!=accSum){
                        acc.sum+=accSum.sum;
                    }
                }
            }
        }
    }
  5. Depuis le répertoire contenant le fichier pom.xml, exécutez la commande suivante :

    mvn package -Dcheckstyle.skip

    L'UDAF est correctement empaquetée lorsque le fichier ASI_UDX-1.0-SNAPSHOT.jar apparaît dans le dossier \ASI_UDX-main\target\.

Utilisation d'une UDAF

Deux méthodes permettent d'utiliser une UDAF dans les déploiements SQL. Le tableau ci-dessous résume les principales différences :

Aspect Méthode 1 : UDAF enregistrée Méthode 2 : JAR au niveau du déploiement
Portée Disponible pour plusieurs déploiements Limitée à un seul déploiement
Mode d'enregistrement Enregistrement via la page de gestion des UDF Téléchargement du JAR sous Additional dependency files dans le déploiement
Appel en SQL Appel par nom enregistré (sans nécessiter CREATE TEMPORARY FUNCTION) Définition d'un alias avec CREATE TEMPORARY FUNCTION … AS 'fully.qualified.ClassName'
Réutilisabilité Élevée — adaptée à la logique métier partagée Faible — liée à un seul déploiement

Méthode 1 : Utiliser une UDAF enregistrée (recommandée)

Enregistrez l'UDAF une fois, puis réutilisez-la dans plusieurs déploiements. Pour les étapes d'enregistrement, consultez la rubrique Gestion des UDF.

Une fois enregistrée sous le nom ASI_UDAF$MySum, appelez-la directement dans votre code SQL :

CREATE TEMPORARY TABLE ASI_UDAF_Source (
  a BIGINT NOT NULL
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE ASI_UDAF_Sink (
  sum  BIGINT
) WITH (
  'connector' = 'print'
);

INSERT INTO ASI_UDAF_Sink
SELECT `ASI_UDAF$MySum`(a)
FROM ASI_UDAF_Source;

Méthode 2 : Télécharger un JAR pour un déploiement spécifique

Sur la page Flink Data Studio > ETL, téléchargez le package JAR en utilisant l'option Additional dependency files située sous More configurations. Définissez ensuite un alias de fonction temporaire dans le code SQL du job.

Le JAR est limité à ce déploiement uniquement et ne peut pas être partagé avec d'autres déploiements.

Si la fonction temporaire est nommée mysum :

CREATE TEMPORARY TABLE ASI_UDAF_Source (
  a   BIGINT
) WITH (
  'connector' = 'datagen'
);

CREATE TEMPORARY TABLE ASI_UDAF_Sink (
  sum  BIGINT
) WITH (
  'connector' = 'print'
);

CREATE TEMPORARY FUNCTION `mysum` AS 'ASI_UDAF.ASI_UDAF$MySum'; -- Create the temporary function mysum.

INSERT INTO ASI_UDAF_Sink
SELECT `mysum`(a)
FROM ASI_UDAF_Source;

Exécution du job

Après avoir développé et déployé le job SQL, accédez à Operation Center > Job O&M. Localisez le job cible et cliquez sur Start dans la colonne Operation.

Une fois le job démarré, l'UDAF agrège le champ a de la table ASI_UDAF_Source et écrit la somme cumulative dans la table ASI_UDAF_Sink.

Étapes suivantes

  • Gestion des UDF : enregistrez et gérez les UDAF pour les réutiliser dans plusieurs déploiements.

  • Optimisation de Flink SQL : configurez le mini-batch pour contrôler l'émission des résultats intermédiaires.

  • Vue d'ensemble : trouvez la version d'Apache Flink correspondant à votre version VVR.