All Products
Search
Document Center

DataHub:Create a subscription

Last Updated:Aug 25, 2026

Subscription feature

When you consume data from a DataHub topic, you must manage your own consumption offsets to resume processing after an application failure. This requires saving your progress and ensuring your offset storage service is highly available, which adds complexity to your application. To simplify this, DataHub provides a subscription service that stores your consumption offsets on the server side. With a few simple configuration steps and minimal code, you gain a highly available offset management service that operates transparently to your application. The subscription service also provides flexible offset reset capabilities, which support at-least-once consumption semantics. For example, if you find a processing error that affected data from a specific time period and want to re-consume the data, you can reset the offset to the corresponding time. Your application automatically detects this change and re-processes the data without requiring a restart.

Create a subscription

Ensure your account has permission to create a subscription for a topic in the specified project. For details, see the Permission Control documentation. Follow these steps:

  • Open the Topic page, click + Subscription in the upper-right corner, fill in the subscription details, and click Create.

    • Subscription Application: The name of the application using this subscription.

    • Description: A detailed description of the subscription.

  • Click the search button under Consumption Checkpoint to view the consumption status of all shards.

Usage example

The subscription feature stores offsets. Although the subscription feature is independent of DataHub's read and write functions (see the Java SDK documentation), it is often used with them when you need to store consumption offsets after reading data.

// 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);
        }
    }
}
  • When the application starts for the first time, it begins consuming data from the earliest available record. As the application runs, you can refresh the subscription page in the web console to see the shard's consumption offset advance.

  • If you manually change the offset by using the Reset Checkpoint feature in the web console while the consumer is running, the application automatically detects the change and resumes consumption from the new offset. To do this, the client catches a SubscriptionOffsetResetException and calls the getSubscriptionOffset method to fetch the latest SubscriptionOffset object from the server.

  • Do not use multiple consumer threads or processes to consume the same shard of a subscription simultaneously. This causes the offset to be overwritten by different consumers, leaving the stored offset in an undefined state. In this scenario, the server throws a SubscriptionSessionInvalidException. Catch this exception, exit the application, and check your design for duplicate consumers.