Utilisez Data Integration de DataWorks pour créer automatiquement des partitions lors de la migration de données depuis ApsaraDB RDS vers MaxCompute.
Prérequis
-
Configurez un environnement DataWorks.
Créez un workflow dans DataWorks. Ce tutoriel utilise le mode basique de DataWorks. Pour plus d'informations, consultez la rubrique Créer un workflow.
-
Ajoutez des sources de données.
Ajoutez une source de données MySQL en tant que source. Pour plus d'informations, consultez la rubrique Configurer une source de données MySQL.
Ajoutez une source de données MaxCompute en tant que destination pour recevoir les données RDS. Pour plus d'informations, consultez la rubrique Configurer une source de données MaxCompute.
Créer automatiquement des partitions
Planifiez une tâche quotidienne pour synchroniser les données d'ApsaraDB RDS vers MaxCompute et créer automatiquement des partitions basées sur la date. Pour plus d'informations sur les tâches de synchronisation de données, consultez la rubrique Développement des données et O&M.
Cet exemple utilise le mode simple de DataWorks. Lors de la création d'un espace de travail, l'option Join the public preview of Data Development (Data Studio) est désactivée par défaut ; les espaces de travail en aperçu public ne sont pas compatibles avec cet exemple.
Connectez-vous à la console DataWorks.
-
Créez une table de destination dans MaxCompute.
Dans le volet de navigation de gauche, cliquez sur Workspace.
Dans la colonne Actions de votre espace de travail, cliquez sur Quick Access > Data Development.
Faites un clic droit sur le workflow créé, puis choisissez .
Sur la page Create Table , sélectionnez une instance de moteur, un schéma et un chemin. Saisissez un Table Name et cliquez sur Create.
Sur la page d'édition de la table, cliquez sur l'icône
pour passer en mode DDL.-
Dans la boîte de dialogue DDL, saisissez l'instruction suivante et cliquez sur Generate Table Schema. Dans la boîte de dialogue Confirm operation qui s'affiche, cliquez sur OK.
CREATE TABLE IF NOT EXISTS ods_user_info_d ( uid STRING COMMENT 'User ID', gender STRING COMMENT 'Gender', age_range STRING COMMENT 'Age range', zodiac STRING COMMENT 'Zodiac sign' ) PARTITIONED BY ( dt STRING ); Cliquez sur Submit to Production Environment.
-
Créez un nœud de synchronisation hors ligne.
Accédez à la page d'analyse des données. Faites un clic droit sur le workflow spécifié et choisissez .
Dans la boîte de dialogue Create Node , saisissez le Name et cliquez sur Confirm.
-
Sélectionnez la source de données, le groupe de ressources et la destination, puis testez la connexion.
Source : la source de données MySQL que vous avez créée.
Resource Group : sélectionnez un groupe de ressources exclusif pour Data Integration.
Destination : la source de données MaxCompute que vous avez créée.
Cliquez sur Next pour configurer la tâche. Pour la source, sélectionnez la table ods_user_info_d et définissez la clé de fragmentation sur
uid. Pour la destination, sélectionnez la table ods_user_info_d et saisissezdt=${bizdate}pour les informations de partition. Définissez le groupe de ressources Tunnel sur public transport resource et le schéma sur default. Pour le mode d'écriture, sélectionnez Insert Overwrite (Clear existing data before writing) et, pour l'écriture des chaînes vides comme null, sélectionnez No.
-
Configurez les paramètres de planification.
Dans le volet de navigation de droite, cliquez sur Properties.
-
Dans la section Scheduling Parameter , le paramètre par défaut est
${bizdate}, qui utilise le format yyyymmdd.RemarqueLa valeur de paramètre par défaut correspond à la Destination configurée pour les Partition Information . Le système remplace automatiquement la valeur de partition par la date métier, qui correspond au jour précédant l'exécution de la tâche, car les jobs ETL traitent généralement les données de la veille. Pour utiliser la date d'exécution actuelle, personnalisez le paramètre.
Vous pouvez personnaliser les formats de date avec les expressions suivantes :
Dans N années :
$[add_months(yyyymmdd,12*N)]Il y a N années :
$[add_months(yyyymmdd,-12*N)]Il y a N mois :
$[add_months(yyyymmdd,-N)]Dans N semaines :
$[yyyymmdd+7*N]Dans N mois :
$[add_months(yyyymmdd,N)]Il y a N semaines :
$[yyyymmdd-7*N]Dans N jours :
$[yyyymmdd+N]Il y a N jours :
$[yyyymmdd-N]Dans N heures :
$[hh24miss+N/24]Il y a N heures :
$[hh24miss-N/24]Dans N minutes :
$[hh24miss+N/24/60]Il y a N minutes :
$[hh24miss-N/24/60]
RemarqueUtilisez des crochets (
[]) pour définir la formule de calcul d'une variable personnalisée. Par exemple :key1=$[yyyy-mm-dd].Par défaut, l'unité de calcul pour les variables personnalisées est le jour. Par exemple,
$[hh24miss-N/24/60]représente le résultat de(yyyymmddhh24miss - (N/24/60 * 1 day)), formaté en hh24miss.L'unité de calcul pour la fonction
add_monthsest le mois. Par exemple,$[add_months(yyyymmdd,12*N)-M/24/60]représente le résultat de(yyyymmddhh24miss - (12 * N * 1 month) - (M/24/60 * 1 day)), formaté enyyyymmdd.
Cliquez sur l'icône
pour exécuter le code.Consultez les résultats dans Runtime Log .
Rétrocharger les données historiques
Pour synchroniser les données historiques antérieures à la tâche planifiée, utilisez la fonctionnalité Complement data dans le Operation Center de DataWorks afin de synchroniser les données et de créer automatiquement des partitions.
-
Filtrez les données historiques de la source ApsaraDB RDS par date.
Dans la section Source du nœud de synchronisation, définissez la condition de Data Filter sur
${bizdate}. Effectuez une rétrocharge de données. Pour plus d'informations, consultez la rubrique Gérer les instances de rétrocharge de données.
-
Dans le journal d'exécution, vérifiez les résultats d'extraction des données ApsaraDB RDS.
Le journal d'exécution indique que MaxCompute a automatiquement créé la partition. La clause
wheredans la configuration du lecteur (where=[20180913]) correspond à lapartitiondans la configuration de l'écritureur (partition=[dt=20180913]Alibaba DI Console, Build 201805310000 . Copyright 2018 Alibaba Group, All rights reserved . Start Job[16961870], traceId [283789484710656#79023#None#None#228255635341196741#None#None#rds_sync], running in Pipeline[basecommon_ 89484710656] The Job[16961870] will run in PhysicsPipeline [basecommon_group_283789484710656_oxs] with requestId [4f44180d-300c-47c3-8ea3-805d2 2018-12-02 03:31:25 : --- Reader: mysql column=["uid","gender","age_range","zodiac"] connection=[{"datasource":"xxx","table":["`ods_user_info_d`"]}]] where=[20180913 ] splitPk=[uid ] Writer: odps isCompress=[false ] partition=[dt=20180913 ] truncate=[true ] datasource=[odps_first ] column=["uid","gender","age_range","zodiac"] emptyAsNull=[false ] table=[ods_user_info_d ] Setting: errorLimit=[{"record":""} ] speed=[{"concurrent":1,"dmu":1,"mbps":"10","throttle":true}] 2018-12-02 03:31:26 : State: 1(SUBMIT) | Total: 0R 0B | Speed: 0R/s 0B/s | Error: 0R 0B | Stage: 0.0% 2018-12-02 03:31:36 : State: 3(RUN) | Total: 0R 0B | Speed: 0R/s 0B/s | Error: 0R 0B | Stage: 0.0% -
Vérifiez le résultat. Exécutez la commande suivante dans le client MaxCompute pour vérifier si les données ont été écrites.
SELECT count(*) from ods_user_info_d where dt = 20180913;
Partitionner des champs non temporels à l'aide d'un hachage
Si vous partitionnez les données selon un champ non temporel tel que la province, Data Integration ne peut pas effectuer de partitionnement automatique. Vous pouvez plutôt hacher un champ issu de RDS pour stocker les enregistrements ayant la même valeur de champ dans la partition MaxCompute correspondante.
-
Synchronisez toutes les données vers une table temporaire dans MaxCompute et créez un nœud de script SQL. Exécutez la commande suivante.
drop table if exists ods_user_t; CREATE TABLE ods_user_t ( dt STRING, uid STRING, gender STRING, age_range STRING, zodiac STRING); --Store the data from the MaxCompute table in the temporary table. insert overwrite table ods_user_t select dt,uid,gender,age_range,zodiac from ods_user_info_d; Créez un nœud de tâche de synchronisation nommé
mysql_to_odps. Il s'agit d'une tâche simple permettant de synchroniser toutes les données de RDS vers MaxCompute. Aucune partition n'est requise.-
Utilisez une instruction SQL pour effectuer un partitionnement dynamique vers la table cible. La commande est la suivante.
drop table if exists ods_user_d; //Create a partitioned ODPS table. This is the final destination table. CREATE TABLE ods_user_d ( uid STRING, gender STRING, age_range STRING, zodiac STRING ) PARTITIONED BY ( dt STRING ); //Run the dynamic partitioning SQL statement. It automatically partitions data based on the dt field of the temporary table. Records with the same value in the dt field are placed in a partition created for that value. //For example, if some records have the value 20181025 in the dt field, a partition dt=20181025 is automatically created in the ODPS partitioned table. //The dynamic partitioning SQL is as follows. //Note that the dt field is included in the SELECT statement. This specifies that partitions are automatically created based on this field. insert overwrite table ods_user_d partition(dt)select dt,uid,gender,age_range,zodiac from ods_user_t; //After the import is complete, delete the temporary table to save storage costs. drop table if exists ods_user_t;Dans MaxCompute, vous pouvez utiliser des instructions SQL pour synchroniser les données.
-
Configurez les trois nœuds dans un workflow pour qu'ils s'exécutent séquentiellement.
L'ordre d'exécution est temporary table, mysql_to_odps et target table.
-
Consultez le processus d'exécution. Concentrez-vous sur le processus de partitionnement dynamique du dernier nœud.
=20181203065434115g3a2eqsa Log view: http://logview.odps.aliyun.com/logview/?h=http://service.odps.aliyun.com/api&p=DataWorks_DOC&i=20181203065434115g3a2eqsa&token=V0NaNDFxUmpn... Summary: resource cost: cpu 0.00 Core * Min, memory 0.00 GB * Min inputs: dataworks_doc.ods_user_t: 20028 (119496 bytes) outputs: dataworks_doc.ods_use xxx 20028 (119176 bytes) Job run time: 0.000 Job run mode: service job Job run engine: execution engine M1: instance count: 1 run time: 0.000 instance time: min: 0.000, max: 0.000, avg: 0.000 input records: ... -
Vérifiez le résultat. Exécutez la commande suivante dans le client MaxCompute pour vérifier les données qui ont été écrites.
SELECT count(*) from ods_user_d where dt = 20180913;