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.
-
Créer une instance ApsaraDB RDS for MySQL.
RemarqueVotre 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.
-
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.
-
Se connecter à l'instance ApsaraDB RDS for MySQL via Data Management (DMS), créez les tables
electric_infoetelectric_info_SortListAggdans la base de donnéeselectric, 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
-
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.
-
L'exemple de code
ASI_UDAFfusionne 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) { } } } -
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).
Connectez-vous à la console Realtime Compute for Apache Flink.
Repérez l'espace de travail cible et cliquez sur Console dans la colonne Actions.
Dans le volet de navigation de gauche, choisissez .
Sous l'onglet UDFs, cliquez sur Register UDF Artifact.
-
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.
RemarqueLe 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.
-
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
Sur la page , cliquez sur New.
Cliquez sur Blank Stream Draft.
Cliquez sur Next.
-
Dans la boîte de dialogue New Draft, configurez les paramètres du job.
Paramètre
Description
Name
Un nom unique pour le job.
RemarqueLe 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
à 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.
-
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_pwpour 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_infoouelectric_info_sortlistagg.port
Numéro de port du service de base de données MySQL.
Aucune.
(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.
Cliquez sur Deploy, puis sur Confirm.
Sur la page , 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
Pour obtenir des informations sur les fonctions intégrées prises en charge par Flink, consultez la section Fonctions intégrées.
Pour en savoir plus sur le déploiement et le démarrage des jobs, consultez les sections Déployer un job et Démarrer un job.
Pour modifier les paramètres d'exécution des jobs, consultez la section Configurer le déploiement de job. Certains paramètres peuvent être mis à jour dynamiquement afin de réduire les temps d'arrêt causés par l'arrêt et le redémarrage des jobs. Pour plus d'informations, consultez la section Mise à l'échelle dynamique et mises à jour des paramètres.
Pour obtenir des informations sur l'utilisation de fonctions personnalisées Python dans les jobs SQL, consultez la section Fonctions personnalisées.