Tous les produits
Search
Centre de documentation

DataHub:Créer un abonnement

Dernière mise à jour :Aug 25, 2026

Fonctionnalité d'abonnement

Lorsque vous consommez des données depuis un topic DataHub, vous devez gérer vos propres offsets de consommation pour reprendre le traitement après une panne applicative. Cette approche implique de sauvegarder votre progression et de garantir la haute disponibilité de votre service de stockage d'offsets, ce qui complexifie votre application. Pour simplifier cette tâche, DataHub propose un service d'abonnement qui stocke vos offsets de consommation côté serveur. Quelques étapes de configuration et un minimum de code suffisent pour bénéficier d'un service de gestion d'offsets hautement disponible, transparent pour votre application. Le service d'abonnement offre également des fonctionnalités flexibles de réinitialisation d'offset, prenant en charge la sémantique de consommation « au moins une fois ». Par exemple, si vous détectez une erreur de traitement ayant affecté les données d'une période spécifique et souhaitez les reconsommer, réinitialisez l'offset à l'heure correspondante. Votre application détecte automatiquement ce changement et retraite les données sans redémarrage.

Créer un abonnement

Assurez-vous que votre compte dispose des autorisations nécessaires pour créer un abonnement pour un topic dans le projet spécifié. Pour plus de détails, consultez la documentation sur le contrôle des autorisations. Suivez les étapes ci-dessous :

  • Ouvrez la page Topic, cliquez sur + Subscription dans le coin supérieur droit, renseignez les détails de l'abonnement et cliquez sur Create.

    • Subscription Application : nom de l'application utilisant cet abonnement.

    • Description : description détaillée de l'abonnement.

  • Cliquez sur le bouton de recherche sous Consumption Checkpoint pour afficher l'état de consommation de tous les shards.

Exemple d'utilisation

La fonctionnalité d'abonnement permet de stocker les offsets. Bien qu'elle soit indépendante des fonctions de lecture et d'écriture de DataHub (consultez la documentation du SDK Java), elle est souvent utilisée conjointement avec celles-ci lorsque vous devez conserver les offsets de consommation après la lecture des données.

// Example of consuming data and committing offsets during the process.
public void offset_consumption(int maxRetry) {
    String endpoint = "<YourEndPoint>";
    String accessId = "<YourAccessId>";
    String accessKey = "<YourAccessKey>";
    String projectName = "<YourProjectName>";
    String topicName = "<YourTopicName>";
    String subId = "<YourSubId>";
    String shardId = "0";
    List<String> shardIds = Arrays.asList(shardId);
    // Create a DatahubClient instance.
    DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
            .setDatahubConfig(
                    new DatahubConfig(endpoint,
                            // Whether to enable binary transfer. This feature is supported by the server since version 2.12.
                            new AliyunAccount(accessId, accessKey), true))
            .build();
    RecordSchema schema = datahubClient.getTopic(projectName, topicName).getRecordSchema();
    OpenSubscriptionSessionResult openSubscriptionSessionResult = datahubClient.openSubscriptionSession(projectName, topicName, subId, shardIds);
    SubscriptionOffset subscriptionOffset = openSubscriptionSessionResult.getOffsets().get(shardId);
    // 1. Get the cursor for the current offset. If the current offset has expired or has never been consumed, get the cursor for the first record within the lifecycle.
    String cursor = "";
    // A sequence number less than 0 indicates that the shard has not been consumed.
    if (subscriptionOffset.getSequence() < 0) {
        // Get the cursor for the first record within the lifecycle.
        cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
    } else {
        // Get the cursor for the next record.
        long nextSequence = subscriptionOffset.getSequence() + 1;
        try {
            // Getting a cursor using SEQUENCE may throw a SeekOutOfRangeException, which indicates that the data at the current cursor has expired.
            cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
        } catch (SeekOutOfRangeException e) {
            // Get the cursor for the first record within the lifecycle.
            cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
        }
    }
    // 2. Read records and save the offset. This example demonstrates reading tuple data and committing the offset every 1,000 records.
    long recordCount = 0L;
    // Read 1,000 records at a time.
    int fetchNum = 1000;
    int retryNum = 0;
    int commitNum = 1000;
    while (retryNum < maxRetry) {
        try {
            GetRecordsResult getRecordsResult = datahubClient.getRecords(projectName, topicName, shardId, schema, cursor, fetchNum);
            if (getRecordsResult.getRecordCount() <= 0) {
                // No data. Sleep and try again.
                System.out.println("no data, sleep 1 second");
                Thread.sleep(1000);
                continue;
            }
            for (RecordEntry recordEntry : getRecordsResult.getRecords()) {
                // Process the data.
                TupleRecordData data = (TupleRecordData) recordEntry.getRecordData();
                System.out.println("field1:" + data.getField("field1") + "\t"
                        + "field2:" + data.getField("field2"));
                // After processing the data, update the offset.
                recordCount++;
                subscriptionOffset.setSequence(recordEntry.getSequence());
                subscriptionOffset.setTimestamp(recordEntry.getSystemTime());
                // Commit the offset every 1000 records.
                if (recordCount % commitNum == 0) {
                    // Commit the offset.
                    Map<String, SubscriptionOffset> offsetMap = new HashMap<>();
                    offsetMap.put(shardId, subscriptionOffset);
                    datahubClient.commitSubscriptionOffset(projectName, topicName, subId, offsetMap);
                    System.out.println("commit offset successful");
                }
            }
            cursor = getRecordsResult.getNextCursor();
        } catch (SubscriptionOfflineException | SubscriptionSessionInvalidException e) {
            // Exit. SubscriptionOfflineException: The subscription is offline. SubscriptionSessionInvalidException: Another client is consuming the same subscription.
            e.printStackTrace();
            throw e;
        } catch (SubscriptionOffsetResetException e) {
            // The offset was reset. You need to get the latest version of the SubscriptionOffset.
            SubscriptionOffset offset = datahubClient.getSubscriptionOffset(projectName, topicName, subId, shardIds).getOffsets().get(shardId);
            subscriptionOffset.setVersionId(offset.getVersionId());
            // After an offset is reset, you must get the new cursor. The method you use to get the cursor should match how the offset was reset.
            // If both sequence and timestamp were set during the reset, you can get the cursor by using either SEQUENCE or SYSTEM_TIME.
            // If only the sequence was set, you must use SEQUENCE.
            // If only the timestamp was set, you must use SYSTEM_TIME.
            // As a general rule, try to get the cursor using SEQUENCE first, then SYSTEM_TIME. If both fail, use OLDEST.
            cursor = null;
            if (cursor == null) {
                try {
                    long nextSequence = offset.getSequence() + 1;
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by SEQUENCE failed, try to get cursor by SYSTEM_TIME");
                }
            }
            if (cursor == null) {
                try {
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SYSTEM_TIME, offset.getTimestamp()).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by SYSTEM_TIME failed, try to get cursor by OLDEST");
                }
            }
            if (cursor == null) {
                try {
                    cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
                    System.out.println("get cursor successful");
                } catch (DatahubClientException exception) {
                    System.out.println("get cursor by OLDEST failed");
                    System.out.println("get cursor failed!!");
                    throw e;
                }
            }
        } catch (LimitExceededException e) {
            // Limit exceeded, retry.
            e.printStackTrace();
            retryNum++;
        } catch (DatahubClientException e) {
            // Other error, retry.
            e.printStackTrace();
            retryNum++;
        } catch (Exception e) {
            e.printStackTrace();
            System.exit(-1);
        }
    }
}
  • Lors du premier démarrage, l'application commence à consommer les données à partir de l'enregistrement disponible le plus ancien. Pendant l'exécution, actualisez la page d'abonnement dans la console web pour voir l'offset de consommation du shard avancer.

  • Si vous modifiez manuellement l'offset via la fonctionnalité Reset Checkpoint dans la console web pendant que le consommateur est en cours d'exécution, l'application détecte automatiquement le changement et reprend la consommation à partir du nouvel offset. Pour ce faire, le client intercepte l'exception SubscriptionOffsetResetException et appelle la méthode getSubscriptionOffset afin de récupérer le dernier objet SubscriptionOffset depuis le serveur.

  • N'utilisez pas plusieurs threads ou processus consommateurs pour consommer simultanément le même shard d'un abonnement. Cela entraînerait l'écrasement de l'offset par différents consommateurs, laissant l'offset stocké dans un état indéfini. Dans ce scénario, le serveur lève l'exception SubscriptionSessionInvalidException. Interceptez cette exception, arrêtez l'application et vérifiez votre conception pour identifier les consommateurs dupliqués.