Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Trier et agréger des données avec une UDAF

Dernière mise à jour :Aug 09, 2026

Cette rubrique utilise des données de terminaux du réseau électrique résidentiel pour illustrer l'utilisation d'une fonction d'agrégation définie par l'utilisateur (UDAF) afin de fusionner et de trier des données dans la console Realtime Compute for Apache Flink.

Exemple de données

La table electric_info contient les données provenant des terminaux du réseau électrique résidentiel. Elle inclut l'ID d'événement (event_id), l'ID utilisateur (user_id), l'heure de l'événement (event_time) et l'état du terminal (status). Vous allez agréger les valeurs de status pour chaque utilisateur et les trier selon event_time.

  • electric_info

    event_id

    user_id

    event_time

    status

    1

    1222

    2023-06-30 11:14:00

    LD

    2

    1333

    2023-06-30 11:12:00

    LD

    3

    1222

    2023-06-30 11:11:00

    TD

    4

    1333

    2023-06-30 11:12:00

    LD

    5

    1222

    2023-06-30 11:15:00

    TD

    6

    1333

    2023-06-30 11:18:00

    LD

    7

    1222

    2023-06-30 11:19:00

    TD

    8

    1333

    2023-06-30 11:10:00

    TD

    9

    1555

    2023-06-30 11:16:00

    TD

    10

    1555

    2023-06-30 11:17:00

    LD

  • Résultat attendu

    user_id

    status

    1222

    TD,LD,TD,TD

    1333

    TD,LD,LD,LD

    1555

    TD,LD

Étape 1 : Préparer la source de données

Cet exemple utilise ApsaraDB RDS comme source de données.

  1. Créer une instance ApsaraDB RDS for MySQL.

    Remarque

    Votre instance ApsaraDB RDS for MySQL doit se trouver dans le même VPC que votre espace de travail Realtime Compute for Apache Flink. Si elles résident dans des VPC distincts, consultez la section Connectivité réseau.

  2. Créer une base de données et un compte.

    Créez une base de données nommée electric ainsi qu'un compte disposant des autorisations de lecture et d'écriture sur cette base de données.

  3. Se connecter à l'instance ApsaraDB RDS for MySQL via Data Management (DMS), créez les tables electric_info et electric_info_SortListAgg dans la base de données electric, puis insérez les données.

    CREATE TABLE `electric_info` (
      event_id bigint NOT NULL PRIMARY KEY COMMENT 'Event ID',
      user_id bigint NOT NULL COMMENT 'User ID', 
      event_time timestamp NOT NULL COMMENT 'Event time',
      status varchar(10) NOT NULL COMMENT 'User terminal status'
    );
    CREATE TABLE `electric_info_SortListAgg` (
      user_id bigint NOT NULL PRIMARY KEY COMMENT 'User ID', 
      status_sort varchar(50) NULL COMMENT 'User terminal status sorted in ascending order by event time'
    );
    -- Prepare data
    INSERT INTO electric_info VALUES 
    (1,1222,'2023-06-30 11:14','LD'),
    (2,1333,'2023-06-30 11:12','LD'),
    (3,1222,'2023-06-30 11:11','TD'),
    (4,1333,'2023-06-30 11:12','LD'),
    (5,1222,'2023-06-30 11:15','TD'),
    (6,1333,'2023-06-30 11:18','LD'),
    (7,1222,'2023-06-30 11:19','TD'),
    (8,1333,'2023-06-30 11:10','TD'),
    (9,1555,'2023-06-30 11:16','TD'),
    (10,1555,'2023-06-30 11:17','LD');

Étape 2 : Enregistrer la UDAF

  1. Téléchargez le package ASI_UDX-1.0-SNAPSHOT.jar.

    Le fichier pom.xml est configuré avec les dépendances minimales requises pour cette fonction personnalisée dans la version 1.17.1 de Flink. Pour plus d'informations sur les fonctions personnalisées, consultez la section Fonctions personnalisées.

  2. L'exemple de code ASI_UDAF fusionne plusieurs lignes en une seule et trie les données selon une colonne spécifiée. Vous pouvez adapter ce code à vos besoins métier.

    package ASI_UDAF;
    import org.apache.commons.lang3.StringUtils;
    import org.apache.flink.table.functions.AggregateFunction;
    import java.util.ArrayList;
    import java.util.Comparator;
    import java.util.Iterator;
    import java.util.List;
    public class ASI_UDAF{
    	/**Accumulator class*/
    	public static class AcList {
    		public  List<String> list;
    	}
    	/**Aggregate function class*/
    	public static class SortListAgg extends AggregateFunction<String,AcList> {
    		public String getValue(AcList asc) {
    			/**Sort the data in the list based on a specific rule*/
    			asc.list.sort(new Comparator<String>() {
    				@Override
    				public int compare(String o1, String o2) {
    					return Integer.parseInt(o1.split("#")[1]) - Integer.parseInt(o2.split("#")[1]);
    				}
    			});
    			/**Traverse the sorted list, extract the required fields, and join them into a string*/
    			List<String> ret = new ArrayList<String>();
    			Iterator<String> strlist = asc.list.iterator();
    			while (strlist.hasNext()) {
    				ret.add(strlist.next().split("#")[0]);
    			}
    			String str = StringUtils.join(ret, ',');
    			return str;
    		}
    		/**Method to create an accumulator*/
    		public AcList createAccumulator() {
    			AcList ac = new AcList();
    			List<String> list = new ArrayList<String>();
    			ac.list = list;
    			return ac;
    		}
    		/**Accumulation method: Add the input data to the accumulator*/
    		public void accumulate(AcList acc, String tuple1) {
    			acc.list.add(tuple1);
    		}
    		/**Retraction method*/
    		public void retract(AcList acc, String num) {
    		}
    	}
    }
  3. Enregistrez la UDAF.

    L'enregistrement d'une UDAF permet de réutiliser son code dans d'autres jobs. Pour les UDAF Java, vous pouvez également télécharger le JAR en tant que fichier de dépendance. Pour en savoir plus, consultez la section Fonctions d'agrégation définies par l'utilisateur (UDAF).

    1. Connectez-vous à la console Realtime Compute for Apache Flink.

    2. Repérez l'espace de travail cible et cliquez sur Console dans la colonne Actions.

    3. Dans le volet de navigation de gauche, choisissez Development > ETL.

    4. Sous l'onglet UDFs, cliquez sur Register UDF Artifact.

  4. Dans la section Click to select, téléchargez le fichier JAR de l'étape 1, puis cliquez sur OK.

    La boîte de dialogue propose deux méthodes d'enregistrement : Upload File et External URL. Vous devez également spécifier un UDF Name et avez la possibilité de télécharger un dependency file.

    Remarque
    • Le fichier JAR de votre UDAF est téléchargé dans le répertoire sql-artifacts du compartiment OSS associé à l'espace de travail.

    • La console Realtime Compute for Apache Flink analyse votre fichier JAR UDAF et détecte les classes utilisant les interfaces Flink UDF, UDAF et UDTF. Elle extrait automatiquement les noms de classe et renseigne le champ Function Name.

  5. Dans la boîte de dialogue Manage Functions, cliquez sur Create Functions.

    La UDF enregistrée apparaît dans la liste UDFs située à gauche de la page de l'éditeur SQL.

Étape 3 : Créer un job Flink

  1. Sur la page Development > ETL, cliquez sur New.

  2. Cliquez sur Blank Stream Draft.

  3. Cliquez sur Next.

  4. Dans la boîte de dialogue New Draft, configurez les paramètres du job.

    Paramètre

    Description

    Name

    Un nom unique pour le job.

    Remarque

    Le nom du job doit être unique au sein du projet actuel.

    Location

    L'emplacement de stockage du job.

    Vous pouvez également cliquer sur l'icône Créer un dossier à côté d'un dossier existant pour créer un sous-dossier.

    Engine version

    La version du moteur Flink pour le job. Elle doit correspondre à la version spécifiée dans votre fichier pom.xml.

    Pour plus de détails sur les versions du moteur, les mappages de versions et les informations sur le cycle de vie, consultez la section Versions du moteur.

  5. Rédigez les instructions DDL et DML.

    -- Create the temporary table electric_info.
    CREATE TEMPORARY TABLE electric_info (
      event_id bigint not null,
      `user_id` bigint not null, 
      event_time timestamp(6) not null,
      status string not null,
      primary key(event_id) not enforced
    ) WITH (
      'connector' = 'mysql',
      'hostname' = 'rm-bp1s1xgll21******.mysql.rds.aliyuncs.com',
      'port' = '3306',
      'username' = 'your_username',
      'password' = '${secret_values.mysql_pw}',
      'database-name' = 'electric',
      'table-name' = 'electric_info'
    );
    CREATE TEMPORARY TABLE electric_info_sortlistagg (
      `user_id` bigint not null, 
      status_sort varchar(50) not null,
      primary key(user_id) not enforced
    ) WITH (
      'connector' = 'mysql',
      'hostname' = 'rm-bp1s1xgll21******.mysql.rds.aliyuncs.com',
      'port' = '3306',
      'username' = 'your_username',
      'password' = '${secret_values.mysql_pw}',
      'database-name' = 'electric',
      'table-name' = 'electric_info_sortlistagg'
    );
    -- Aggregate data from the electric_info table and insert it into the electric_info_sortlistagg table.
    -- Pass a concatenated string of status and event_time as a parameter to the registered custom function ASI_UDAF$SortListAgg.
    INSERT INTO electric_info_sortlistagg 
    SELECT `user_id`, `ASI_UDAF$SortListAgg`(CONCAT(status,'#',CAST(UNIX_TIMESTAMP(event_time) as STRING)))
    FROM electric_info GROUP BY user_id;

    Le tableau suivant décrit les paramètres. Modifiez-les selon vos besoins réels. Pour plus d'informations sur les paramètres du connecteur MySQL, consultez la section Connecteur MySQL.

    Paramètre

    Description

    Notes

    connector

    Type de connecteur.

    Dans cet exemple, la valeur est fixée à mysql.

    hostname

    Adresse IP ou nom d'hôte de la base de données MySQL.

    Cet exemple utilise le point de terminaison interne de l'instance ApsaraDB RDS for MySQL.

    username

    Nom d'utilisateur du service de base de données MySQL.

    Aucune.

    password

    Mot de passe du service de base de données MySQL.

    Cet exemple utilise une variable nommée mysql_pw pour le mot de passe afin d'éviter les risques de sécurité. Pour plus d'informations, consultez la section Variables.

    database-name

    Nom de la base de données MySQL.

    Cet exemple utilise electric, la base de données créée à l'Étape 1 : Préparer la source de données.

    table-name

    Nom de la table MySQL.

    Dans cet exemple, définissez cette valeur sur electric_info ou electric_info_sortlistagg.

    port

    Numéro de port du service de base de données MySQL.

    Aucune.

  6. (Facultatif) Dans le coin supérieur droit, cliquez sur Validate et Debug. Pour plus d'informations sur ces fonctionnalités, consultez la section Présentation du développement de jobs.

  7. Cliquez sur Deploy, puis sur Confirm.

  8. Sur la page O&M > Deployments, repérez le job cible, cliquez sur Start dans la colonne Actions, puis sélectionnez Initial Mode.

Étape 4 : Consulter le résultat

Dans ApsaraDB RDS, exécutez l'instruction suivante pour afficher les résultats agrégés et triés.

SELECT * FROM `electric_info_sortlistagg`;

La sortie confirme que les états de chaque utilisateur ont été correctement agrégés et triés : user_id=1222 correspond à status_sort=TD,LD,TD,TD ; user_id=1333 correspond à status_sort=TD,LD,LD,LD ; et user_id=1555 correspond à status_sort=TD,LD.

Références