Chargez les données vers MaxCompute en parallèle en utilisant l'interface TableTunnel avec le service ExecutorService de Java. Chaque thread possède un bloc exclusif et un objet RecordWriter ; les blocs ne doivent pas être partagés entre les threads.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un projet MaxCompute contenant une table partitionnée vers laquelle écrire
Le SDK Java MaxCompute ajouté aux dépendances de votre projet
-
Un ID AccessKey et un secret AccessKey stockés en tant que variables d'environnement :
ALIBABA_CLOUD_ACCESS_KEY_IDALIBABA_CLOUD_ACCESS_KEY_SECRET
(Recommandé) Un utilisateur RAM disposant des autorisations minimales requises, plutôt qu'un compte racine Alibaba Cloud. Les identifiants du compte racine présentent un risque élevé, car la paire AccessKey dispose d'autorisations sur toutes les opérations API. Créez un utilisateur RAM dans la console Resource Access Management (RAM).
Fonctionnement
Le processus de chargement comporte quatre étapes :
Créez une seule session
UploadSessionpour la partition de table cible.Pour chaque thread, ouvrez un objet
RecordWriterdédié en appelantuploadSession.openRecordWriter(blockId). L'ID de bloc identifie de manière unique le segment de données de ce thread ; deux threads ne partagent jamais le même bloc.Chaque thread remplit les enregistrements et les écrit via son propre objet
RecordWriter, puis ferme l'objet d'écriture.Une fois tous les threads terminés, validez la session de chargement avec la liste complète des ID de bloc.
Chargement des données avec plusieurs threads
L'exemple suivant crée 10 threads de chargement, chacun écrivant 10 enregistrements dans son propre bloc.
import java.io.IOException;
import java.util.ArrayList;
import java.util.Date;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import com.aliyun.odps.Column;
import com.aliyun.odps.Odps;
import com.aliyun.odps.PartitionSpec;
import com.aliyun.odps.TableSchema;
import com.aliyun.odps.account.Account;
import com.aliyun.odps.account.AliyunAccount;
import com.aliyun.odps.data.Record;
import com.aliyun.odps.data.RecordWriter;
import com.aliyun.odps.tunnel.TableTunnel;
import com.aliyun.odps.tunnel.TunnelException;
import com.aliyun.odps.tunnel.TableTunnel.UploadSession;
// Each thread owns one block (RecordWriter) and writes 10 records to it.
class UploadThread implements Callable<Boolean> {
private long id;
private RecordWriter recordWriter;
private Record record;
private TableSchema tableSchema;
public UploadThread(long id, RecordWriter recordWriter, Record record,
TableSchema tableSchema) {
this.id = id;
this.recordWriter = recordWriter;
this.record = record;
this.tableSchema = tableSchema;
}
@Override
public Boolean call() {
// Populate each column based on its data type
for (int i = 0; i < tableSchema.getColumns().size(); i++) {
Column column = tableSchema.getColumn(i);
switch (column.getType()) {
case BIGINT:
record.setBigint(i, 1L);
break;
case BOOLEAN:
record.setBoolean(i, true);
break;
case DATETIME:
record.setDatetime(i, new Date());
break;
case DOUBLE:
record.setDouble(i, 0.0);
break;
case STRING:
record.setString(i, "sample");
break;
default:
throw new RuntimeException("Unknown column type: "
+ column.getType());
}
}
boolean success = true;
try {
for (int i = 0; i < 10; i++) {
recordWriter.write(record);
}
} catch (IOException e) {
success = false;
e.printStackTrace();
} finally {
recordWriter.close();
}
return success;
}
}
public class UploadThreadSample {
// Load credentials from environment variables—never hardcode them in source code.
private static String accessId = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
private static String accessKey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
private static String odpsUrl = "<http://service.odps.aliyun.com/api>";
// Set tunnelUrl to a specific Tunnel endpoint, or leave it blank to use the public endpoint.
// The example below uses the classic network Tunnel endpoint in the China (Shanghai) region.
private static String tunnelUrl = "<http://dt.cn-shanghai.maxcompute.aliyun-inc.com>";
private static String project = "<your project>";
private static String table = "<your table name>";
private static String partition = "<your partition spec>";
private static int threadNum = 10;
public static void main(String args[]) {
Account account = new AliyunAccount(accessId, accessKey);
Odps odps = new Odps(account);
odps.setEndpoint(odpsUrl);
odps.setDefaultProject(project);
try {
TableTunnel tunnel = new TableTunnel(odps);
tunnel.setEndpoint(tunnelUrl);
PartitionSpec partitionSpec = new PartitionSpec(partition);
// Create a single upload session for the target partition
UploadSession uploadSession = tunnel.createUploadSession(project,
table, partitionSpec);
System.out.println("Session Status is : " + uploadSession.getStatus().toString());
// Assign one RecordWriter (block) per thread
ExecutorService pool = Executors.newFixedThreadPool(threadNum);
ArrayList<Callable<Boolean>> callers = new ArrayList<Callable<Boolean>>();
for (int i = 0; i < threadNum; i++) {
RecordWriter recordWriter = uploadSession.openRecordWriter(i);
Record record = uploadSession.newRecord();
callers.add(new UploadThread(i, recordWriter, record,
uploadSession.getSchema()));
}
// Run all threads and wait for them to finish
pool.invokeAll(callers);
pool.shutdown();
// Commit all blocks
Long[] blockList = new Long[threadNum];
for (int i = 0; i < threadNum; i++)
blockList[i] = Long.valueOf(i);
uploadSession.commit(blockList);
System.out.println("upload success!");
} catch (TunnelException e) {
e.printStackTrace();
} catch (IOException e) {
e.printStackTrace();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}
Remplacez les espaces réservés suivants avant d'exécuter le code :
| Espace réservé | Description | Exemple |
|---|---|---|
<your project> |
Nom du projet MaxCompute | my_project |
<your table name> |
Nom de la table cible | my_table |
<your partition spec> |
Spécification de partition | ds=20240101 |
Configuration de l'endpoint Tunnel
Définissez tunnelUrl en fonction de votre type de réseau :
| Scénario | tunnelUrl |
|---|---|
| Laisser vide | Endpoint public utilisé par défaut |
| Réseau interne (réseau classique, Chine (Shanghai)) | http://dt.cn-shanghai.maxcompute.aliyun-inc.com |
| Autres régions ou types de réseau | Consultez les Endpoints |
Vérification du chargement
Après l'affichage du message upload success! par le programme, vérifiez les données dans MaxCompute en exécutant la commande suivante :
SELECT COUNT(*) FROM <your table name> WHERE <your partition spec>;
Le nombre obtenu doit correspondre au total des enregistrements écrits par l'ensemble des threads. Dans cet exemple : 10 threads x 10 enregistrements = 100 enregistrements.
Sécurité
Stockez les identifiants AccessKey dans des variables d'environnement, et non dans le code source. Pour les charges de travail de production, utilisez un utilisateur RAM disposant des autorisations minimales requises plutôt que les identifiants du compte racine. Pour créer un utilisateur RAM, accédez à la console Resource Access Management (RAM).
Étapes suivantes
Consultez les Endpoints pour trouver l'endpoint Tunnel correspondant à votre région et à votre type de réseau.