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 instancesTunnelWorkersur un même serveur, nous vous recommandons d'utiliser la même configurationTunnelWorkerConfigpour toutes les instances.L'instance
TunnelWorkernécessite une période de préchauffage pour l'initialisation, définie par le paramètreHeartbeatIntervaldans la configurationTunnelWorkerConfig. La valeur par défaut est 30 secondes.Lorsque le client
TunnelWorkers'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éthodeshutdownenregistrée pour la classeChannelet 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
OTSTunnelExpiredse 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
Une table de données a été créée. Pour plus d'informations, consultez les rubriques Utilisation de la console Tablestore, Utilisation de l'interface CLI Tablestore et Utilisation des SDK Tablestore.
Vous avez obtenu le endpoint de l'instance où réside la table de données. Pour plus d'informations, consultez la rubrique Obtention d'un endpoint d'une instance Tablestore.
Les identifiants d'accès sont configurés. Pour plus d'informations, consultez la rubrique Configuration des identifiants d'accès.
Prise en main du service Tunnel
-
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_IDetTABLESTORE_ACCESS_KEY_SECRETsont configurées. La variable d'environnementTABLESTORE_ACCESS_KEY_IDspécifie l'AccessKey ID de votre compte Alibaba Cloud ou de votre utilisateur RAM. La variable d'environnementTABLESTORE_ACCESS_KEY_SECRETspécifie l'AccessKey Secret de votre compte Alibaba Cloud ou de votre utilisateur RAM.AvertissementUn 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
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.
Un client de tunnel fournit l'opération NewTunnelClientWithToken que vous pouvez appeler pour initialiser une instance
TunnelClientbasée sur des identifiants d'accès temporaires. Cette rubrique fournit un exemple de code pour initialiser une instanceTunnelClientà 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.
-
-
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) -
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()),
)
}