Cette rubrique décrit les classes et méthodes principales de MapReduce dans le SDK Java MaxCompute (odps-sdk-mapred).
Ajout de la dépendance du SDK
Recherchez odps-sdk-mapred dans le référentiel Maven pour identifier la dernière version. Ajoutez la dépendance suivante à votre projet :
<dependency>
<groupId>com.aliyun.odps</groupId>
<artifactId>odps-sdk-mapred</artifactId>
<version>0.40.10-public</version>
</dependency>
Types de données
MaxCompute MapReduce prend en charge les types de données suivants. Le tableau ci-dessous indique leur correspondance avec les types Java.
|
Type MaxCompute |
Type Java |
|
BIGINT |
Long |
|
STRING |
String |
|
DOUBLE |
Double |
|
BOOLEAN |
Boolean |
|
DATETIME |
Date |
|
DECIMAL |
BigDecimal |
Présentation des classes
|
Classe |
Description |
|
|
Classe de base pour les mappeurs définis par l'utilisateur. Elle convertit les enregistrements des tables d'entrée en paires clé-valeur et les transmet au réducteur. Les tâches qui ignorent l'étape de réduction et écrivent directement les résultats sont appelées tâches MapOnly. |
|
|
Classe de base pour les réducteurs définis par l'utilisateur. Elle agrège l'ensemble des valeurs associées à chaque clé. |
|
|
Fournit le contexte d'exécution de la tâche. Cet objet est transmis en tant que paramètre d'entrée aux méthodes de cycle de vie de |
|
|
Soumet et gère les tâches. Il prend en charge la soumission bloquante (synchrone) et non bloquante (asynchrone). |
|
|
Représente une instance de tâche en cours d'exécution. Utilisez cet objet pour suivre l'état, attendre la fin de l'exécution et récupérer les compteurs. |
|
|
Contient la configuration d'une tâche MapReduce. Définissez un objet |
MapperBase
Le framework appelle les méthodes du cycle de vie du mappeur dans l'ordre suivant : setup une fois au début, map une fois par enregistrement d'entrée, et cleanup une fois à la fin.
public class WordCountMapper extends MapperBase {
@Override
public void setup(TaskContext context) throws IOException {
// Initialize resources (for example, load lookup tables)
}
@Override
public void map(long key, Record record, TaskContext context) throws IOException {
// Process each input record and emit key-value pairs.
// MaxCompute uses (long key, Record record), not (KEYIN key, VALUEIN value).
Record mapKey = context.createMapOutputKeyRecord();
Record mapValue = context.createMapOutputValueRecord();
// Populate mapKey and mapValue, then emit:
context.write(mapKey, mapValue);
}
@Override
public void cleanup(TaskContext context) throws IOException {
// Release resources or emit final accumulated results
}
}
|
Méthode |
Description |
|
|
Appelée une fois avant le premier appel à |
|
|
Appelée une fois par enregistrement d'entrée. Émettez des paires clé-valeur avec |
|
|
Appelée une fois après le dernier appel à |
ReducerBase
Le framework appelle les méthodes du cycle de vie du réducteur dans l'ordre suivant : setup une fois au début, reduce une fois par groupe de clés uniques, et cleanup une fois à la fin.
public class WordCountReducer extends ReducerBase {
@Override
public void setup(TaskContext context) throws IOException {
// Initialize resources
}
@Override
public void reduce(Record key, Iterator<Record> values, TaskContext context) throws IOException {
// Aggregate all values for this key
long count = 0;
while (values.hasNext()) {
Record val = values.next();
count += val.getBigint(0);
}
Record output = context.createOutputRecord();
output.set(0, key.getString(0));
output.set(1, count);
context.write(output);
}
@Override
public void cleanup(TaskContext context) throws IOException {
// Release resources
}
}
|
Méthode |
Description |
|
|
Appelée une fois avant le premier appel à |
|
|
Appelée une fois par groupe de clés uniques. Toutes les valeurs associées à la clé sont transmises via |
|
|
Appelée une fois après le dernier appel à |
TaskContext
L'objet TaskContext est transmis à chaque méthode du cycle de vie et permet d'accéder aux tables de sortie, aux ressources et aux utilitaires du framework.
|
Méthode |
Description |
|
|
Renvoie les informations relatives aux tables de sortie. |
|
|
Crée un enregistrement pour la table de sortie par défaut. |
|
|
Crée un enregistrement pour la table de sortie identifiée par |
|
|
Crée un enregistrement pour la clé de sortie de l'étape map. |
|
|
Crée un enregistrement pour la valeur de sortie de l'étape map. |
|
|
Écrit un enregistrement dans la table de sortie par défaut. Cette méthode peut être appelée plusieurs fois pendant l'étape de réduction. |
|
|
Écrit un enregistrement dans la table de sortie identifiée par |
|
|
Émet une paire clé-valeur pendant l'étape map. Cette méthode peut être appelée plusieurs fois. |
|
|
Lit une ressource de fichier par son nom. |
|
|
Lit une ressource de table par son nom. |
|
|
Renvoie le compteur portant le nom spécifié. |
|
|
Renvoie le compteur portant le nom spécifié dans le groupe indiqué. |
|
|
Envoie un signal de maintien actif (heartbeat) au framework MapReduce pour éviter le dépassement du délai d'inactivité du worker. |
Délai d'inactivité du worker
Le délai d'inactivité par défaut du worker est de 10 minutes et ne peut pas être modifié. Si un worker n'appelle pas progress() dans un délai de 10 minutes, le framework met fin au processus du worker et la tâche map ou reduce échoue. Appelez progress() périodiquement dans les tâches de longue durée pour maintenir les workers actifs. Cette méthode envoie uniquement un signal de maintien actif ; elle ne signale pas la progression de la tâche.
JobConf
L'objet JobConf contient toute la configuration d'une tâche MapReduce, y compris les classes de mappeur et de réducteur, les schémas clé-valeur et les déclarations de ressources.
|
Méthode |
Description |
|
|
Définit la classe du mappeur pour la tâche. |
|
|
Définit la classe du réducteur pour la tâche. |
|
|
Définit une classe de combiner. Un combiner pré-agrège les enregistrements partageant la même clé lors de l'étape map, ce qui réduit le volume de données transféré vers les réducteurs. |
|
|
Définit le schéma des clés transmises du mappeur au réducteur. |
|
|
Définit le schéma des valeurs transmises du mappeur au réducteur. |
|
|
Définit les colonnes utilisées pour trier les clés avant leur transmission aux réducteurs. |
|
|
Définit les colonnes utilisées pour regrouper les clés. Les colonnes de regroupement doivent constituer un sous-ensemble des colonnes de tri. |
|
|
Définit les colonnes de clé de partition. Par défaut, toutes les colonnes de clé sont utilisées. |
|
|
Déclare les ressources accessibles aux mappeurs et aux réducteurs. Un mappeur ou un réducteur ne peut lire que les ressources déclarées ici. |
|
|
Définit la taille du fragment d'entrée en Mo. Valeur par défaut : 256 Mo. |
|
|
Définit le nombre de tâches de réduction. Valeur par défaut : un quart du nombre de tâches map. |
|
|
Définit la mémoire allouée par worker map en Mo. Valeur par défaut : 2048 Mo. |
|
|
Définit la mémoire allouée par worker reduce en Mo. Valeur par défaut : 2048 Mo. |
Distribution et regroupement des enregistrements lors des étapes map et reduce
Lors de l'étape map, le framework calcule un hachage de chaque enregistrement de sortie en fonction des colonnes de clé de partition afin de déterminer quel réducteur doit le recevoir. Les enregistrements sont triés selon les colonnes de tri avant d'être envoyés.
Lors de l'étape reduce, les enregistrements sont regroupés par les colonnes de regroupement. Tous les enregistrements partageant la même clé de regroupement sont transmis ensemble à un seul appel de reduce().
Les colonnes de regroupement sont sélectionnées parmi les colonnes de tri. Les colonnes de tri et les colonnes de clé de partition doivent exister dans les clés.
JobClient
|
Méthode |
Description |
|
|
Soumet une tâche en mode bloquant (synchrone). L'appel bloque jusqu'à la fin de la tâche et renvoie un objet |
|
|
Soumet une tâche en mode non bloquant (asynchrone). L'appel renvoie immédiatement un objet |
RunningJob
|
Méthode |
Description |
|
|
Renvoie l'ID de l'instance de tâche. Utilisez cet ID pour consulter les journaux opérationnels et gérer la tâche. |
|
|
Renvoie |
|
|
Renvoie |
|
|
Attend la fin d'une instance de tâche. Cette méthode est utilisée pour les tâches soumises en mode synchrone. |
|
|
Renvoie l'état actuel de l'instance de tâche. |
|
|
Met fin à la tâche en cours d'exécution. |
|
|
Renvoie toutes les données des compteurs pour la tâche. |
InputUtils
|
Méthode |
Description |
|
|
Ajoute une seule table d'entrée à la tâche. Des appels multiples à cette méthode ajoutent chaque table à la file d'attente d'entrée. |
|
|
Définit plusieurs tables d'entrée en une seule fois. |
OutputUtils
|
Méthode |
Description |
|
|
Ajoute une seule table de sortie à la tâche. Des appels multiples à cette méthode ajoutent chaque table à la file d'attente de sortie. |
|
|
Définit plusieurs tables de sortie en une seule fois. |
Pipeline (modèle MapReduce étendu)
L'objet Pipeline constitue le point d'entrée du modèle MapReduce étendu, qui permet d'enchaîner plusieurs mappeurs et réducteurs au sein d'une seule tâche. Utilisez Pipeline.builder() pour construire le pipeline.
Méthodes du générateur
public Builder addMapper(Class<? extends Mapper> mapper)
public Builder addMapper(Class<? extends Mapper> mapper,
Column[] keySchema, Column[] valueSchema, String[] sortCols,
SortOrder[] order, String[] partCols,
Class<? extends Partitioner> theClass, String[] groupCols)
public Builder addReducer(Class<? extends Reducer> reducer)
public Builder addReducer(Class<? extends Reducer> reducer,
Column[] keySchema, Column[] valueSchema, String[] sortCols,
SortOrder[] order, String[] partCols,
Class<? extends Partitioner> theClass, String[] groupCols)
public Builder setOutputKeySchema(Column[] keySchema)
public Builder setOutputValueSchema(Column[] valueSchema)
public Builder setOutputKeySortColumns(String[] sortCols)
public Builder setOutputKeySortOrder(SortOrder[] order)
public Builder setPartitionColumns(String[] partCols)
public Builder setPartitionerClass(Class<? extends Partitioner> theClass)
public Builder setOutputGroupingColumns(String[] cols)
Exemple
L'exemple suivant enchaîne un mappeur et deux réducteurs (TokenizerMapper → SumReducer → IdentityReducer) :
Job job = new Job();
Pipeline pipeline = Pipeline.builder()
.addMapper(TokenizerMapper.class)
.setOutputKeySchema(
new Column[] { new Column("word", OdpsType.STRING) })
.setOutputValueSchema(
new Column[] { new Column("count", OdpsType.BIGINT) })
.addReducer(SumReducer.class)
.setOutputKeySchema(
new Column[] { new Column("count", OdpsType.BIGINT) })
.setOutputValueSchema(
new Column[] { new Column("word", OdpsType.STRING),
new Column("count", OdpsType.BIGINT) })
.addReducer(IdentityReducer.class).createPipeline();
job.setPipeline(pipeline);
job.addInput(...)
job.addOutput(...)
job.submit();
Pour enchaîner un seul mappeur avec un réducteur, utilisezJobConfau lieu dePipeline. Familiarisez-vous avec l'API MapReduce standard avant d'utiliser le modèle étendu.