Tous les produits
Search
Centre de documentation

Tablestore:Prise en main

Dernière mise à jour :Aug 18, 2026

Le service Tunnel vous permet de consommer les données d'une table. Cette rubrique explique comment prendre en main le service Tunnel à l'aide du SDK Tablestore pour Go. Avant d'utiliser le service Tunnel, assurez-vous de bien connaître les remarques d'utilisation associées.

Remarques d'utilisation

  • Par défaut, le système démarre un pool de threads pour lire et traiter les données selon la configuration TunnelWorkerConfig. Si vous souhaitez démarrer plusieurs instances TunnelWorker sur un même serveur, nous vous recommandons d'utiliser la même configuration TunnelWorkerConfig pour toutes les instances.

  • L'instance TunnelWorker nécessite une période de préchauffage pour l'initialisation, définie par le paramètre HeartbeatInterval dans la configuration TunnelWorkerConfig. La valeur par défaut est 30 secondes.

  • Lorsque le client TunnelWorker s'arrête en raison d'une sortie inattendue ou d'une interruption manuelle, il recycle automatiquement les ressources selon l'une des méthodes suivantes : libération du pool de threads, appel automatique de la méthode shutdown enregistrée pour la classe Channel et arrêt du tunnel.

  • La durée de conservation des journaux incrémentiels dans les tunnels est identique à celle des journaux Stream. Les journaux Stream peuvent être conservés jusqu'à sept jours. Par conséquent, les journaux incrémentiels dans les tunnels peuvent également être conservés jusqu'à sept jours.

  • Si vous créez un tunnel pour consommer des données différentielles ou incrémentielles, tenez compte des points suivants :

    • Lors de la consommation des données complètes, si le tunnel ne parvient pas à terminer la consommation des données complètes dans le délai de conservation des journaux incrémentiels (sept jours au maximum), une erreur OTSTunnelExpired se produit lorsque le tunnel commence à consommer les journaux incrémentiels. Le tunnel ne peut alors plus consommer les journaux incrémentiels.

      Si vous estimez que le tunnel ne pourra pas terminer la consommation des données complètes dans le délai imparti, contactez le support technique Tablestore .

    • Lors de la consommation des données incrémentielles, si le tunnel ne parvient pas à terminer la consommation des journaux incrémentiels dans le délai de conservation (sept jours au maximum), il peut consommer les données à partir des dernières données disponibles. Dans ce cas, certaines données risquent de ne pas être consommées.

  • Une fois qu'un tunnel a expiré, Tablestore peut le désactiver. Si un tunnel reste à l'état désactivé pendant plus de 30 jours, il est supprimé. Vous ne pouvez pas restaurer un tunnel supprimé.

Prérequis

Prise en main du service Tunnel

  1. Initialisez une instance TunnelClient.

    Lors de l'initialisation d'une instance TunnelClient, vous pouvez utiliser des identifiants d'accès à long terme ou des identifiants d'accès temporaires pour l'authentification.

    • Utilisation d'identifiants d'accès à long terme pour l'initialisation

      Assurez-vous que les variables d'environnement TABLESTORE_ACCESS_KEY_ID et TABLESTORE_ACCESS_KEY_SECRET sont configurées. La variable d'environnement TABLESTORE_ACCESS_KEY_ID spécifie l'AccessKey ID de votre compte Alibaba Cloud ou de votre utilisateur RAM. La variable d'environnement TABLESTORE_ACCESS_KEY_SECRET spécifie l'AccessKey Secret de votre compte Alibaba Cloud ou de votre utilisateur RAM.

      Avertissement

      Un compte Alibaba Cloud dispose d'un accès complet à toutes les ressources du compte. La divulgation de la paire de clés AccessKey du compte Alibaba Cloud constitue une menace critique pour le système. Nous vous recommandons donc d'utiliser une paire de clés AccessKey appartenant à un utilisateur RAM disposant des autorisations minimales requises pour initialiser une instance TunnelClient.

      // Set the endpoint parameter to the endpoint of the Tablestore instance. Example: https://instance.cn-hangzhou.ots.aliyuncs.com. 
      // Specify the name of the instance. 
      // Specify the AccessKey ID and AccessKey secret of your Alibaba Cloud account or a RAM user. 
      endpoint := "yourEndpoint"
      instance := "yourInstance"
      accessKeyId := os.Getenv("TABLESTORE_ACCESS_KEY_ID")
      accessKeySecret := os.Getenv("TABLESTORE_ACCESS_KEY_SECRET")
      tunnelClient := tunnel.NewTunnelClient(endpoint, instance, accessKeyId, accessKeySecret)                    
    • Utilisation d'identifiants d'accès temporaires pour l'initialisation

      1. Si vous souhaitez utiliser le SDK Tablestore pour Go afin d'accéder temporairement à Tablestore, vous pouvez utiliser Security Token Service (STS) pour générer des identifiants d'accès temporaires. Pour plus d'informations, consultez la rubrique Configuration des identifiants d'accès temporaires.

      2. Un client de tunnel fournit l'opération NewTunnelClientWithToken que vous pouvez appeler pour initialiser une instance TunnelClient basée sur des identifiants d'accès temporaires. Cette rubrique fournit un exemple de code pour initialiser une instance TunnelClient à l'aide d'identifiants d'accès temporaires pouvant être actualisés périodiquement. Pour plus d'informations, consultez la section Annexe : Exemple de code pour l'initialisation d'une instance TunnelClient à l'aide d'identifiants d'accès temporaires.

  2. Créez un tunnel.

    req := &tunnel.CreateTunnelRequest{
       TableName:  "<TABLE_NAME>",
       TunnelName: "<TUNNEL_NAME>",
       Type:       tunnel.TunnelTypeBaseStream, // Create a BaseAndStream tunnel. 
    }
    resp, err := tunnelClient.CreateTunnel(req)
    if err != nil {
       log.Fatal("create test tunnel failed", err)
    }
    log.Println("tunnel id is", resp.TunnelId)
  3. Spécifiez une fonction de rappel personnalisée pour démarrer la consommation automatique des données.

    // Specify a custom callback function. 
    func exampleConsumeFunction(ctx *tunnel.ChannelContext, records []*tunnel.Record) error {
        fmt.Println("user-defined information", ctx.CustomValue)
        for _, rec := range records {
            fmt.Println("tunnel record detail:", rec.String())
        }
        fmt.Println("a round of records consumption finished")
        return nil
    }
    
    // Configure the callback function. Information about the callback function is passed to SimpleProcessFactory. Configure TunnelWorkerConfig for the consumer. 
    workConfig := &tunnel.TunnelWorkerConfig{
       ProcessorFactory: &tunnel.SimpleProcessFactory{
          CustomValue: "user custom interface{} value",
          ProcessFunc: exampleConsumeFunction,
       },
    }
    
    // Use TunnelDaemon to continuously consume the specified tunnel. 
    tunnelId := "<TUNNEL_ID>"
    daemon := tunnel.NewTunnelDaemon(tunnelClient, tunnelId, workConfig)
    log.Fatal(daemon.Run())

Annexe : Exemple de code pour l'initialisation d'une instance TunnelClient à l'aide d'identifiants d'accès temporaires

import (
    otscommon "github.com/aliyun/aliyun-tablestore-go-sdk/common"
    "github.com/aliyun/aliyun-tablestore-go-sdk/tunnel"
    "sync"
    "time"
)

type RefreshClient struct {
    lastRefresh          time.Time
    refreshIntervalInMin int
}

func NewRefreshClient(intervalInMin int) *RefreshClient {
    return &RefreshClient{
        refreshIntervalInMin: intervalInMin,
    }
}

func (c *RefreshClient) IsExpired() bool {
    now := time.Now()
    if c.lastRefresh.IsZero() || now.Sub(c.lastRefresh) > time.Duration(c.refreshIntervalInMin)*time.Minute {
        return true
    }

    return false
}

func (c *RefreshClient) Update() {
    c.lastRefresh = time.Now()
}

type clientCredentials struct {
    accessKeyID     string
    accessKeySecret string
    securityToken   string
}

func newClientCredentials(accessKeyID string, accessKeySecret string, securityToken string) *clientCredentials {
    return &clientCredentials{accessKeyID: accessKeyID, accessKeySecret: accessKeySecret, securityToken: securityToken}
}

func (c *clientCredentials) GetAccessKeyID() string {
    return c.accessKeyID
}

func (c *clientCredentials) GetAccessKeySecret() string {
    return c.accessKeySecret
}

func (c *clientCredentials) GetSecurityToken() string {
    return c.securityToken
}

type OTSCredentialsProvider struct {
    refresh *RefreshClient
    cred    *clientCredentials
    lock    sync.Mutex
}

func NewOTSCredentialsProvider() *OTSCredentialsProvider {
    return &OTSCredentialsProvider{
        // Modify the refresh cycle for temporary access credentials based on your business requirements. The refresh cycle must be shorter than the validity period of the temporary access credentials. 
        refresh: NewRefreshClient(30),
    }
}

func (p *OTSCredentialsProvider) renewCredentials() error {
    if p.cred == nil || p.refresh.IsExpired() {
        // Obtain temporary access credentials. You can call the AssumeRole operation of RAM to obtain the AccessKey ID, AccessKey secret, security token, and validity period of the temporary access credentials. 
        // Configure the following parameters. For information about RAM SDKs, see the documentation of RAM.  
        // resp, err := GetUserOtsStsToken()
        accessKeyId := ""
        accessKeySecret := ""
        stsToken := ""
        p.cred = newClientCredentials(accessKeyId, accessKeySecret, stsToken)
        p.refresh.Update()
    }

    return nil
}

func (p *OTSCredentialsProvider) GetCredentials() otscommon.Credentials {
    p.lock.Lock()
    defer p.lock.Unlock()

    if err := p.renewCredentials(); err != nil {
        // log error
        if p.cred == nil {
            return newClientCredentials("", "", "")
        }
    }

    return p.cred
}

// NewTunnelClientWithToken is used to initialize a TunnelClient instance with the feature of refreshing temporary access credentials. 
func NewTunnelClientWithToken(endpoint, instanceName, accessId, accessKey, token string) tunnel.TunnelClient {
    return tunnel.NewTunnelClientWithToken(
        endpoint,
        instanceName,
        "",
        "",
        "",
        nil,
        tunnel.SetCredentialsProvider(NewOTSCredentialsProvider()),
    )
}