Le connecteur TSDB for InfluxDB sera rendu obsolète à partir de la version 11.7. Une fois obsolète, il sera retiré de la console et ne bénéficiera plus de mises à jour fonctionnelles ni de maintenance. Pour consulter le calendrier d'obsolescence, reportez-vous à la rubrique Fin de prise en charge du connecteur TSDB for InfluxDB. Migrez vos charges de travail dès que possible pour éviter toute interruption de vos jobs de production.
Le connecteur TSDB for InfluxDB écrit les données en streaming issues d'une table de destination Flink SQL dans une instance TSDB for InfluxDB Ververica Runtime (VVR). TSDB for InfluxDB est une base de données temporelle optimisée pour des débits d'écriture et de requête élevés, couramment utilisée pour la surveillance DevOps, les métriques d'applications et les données de capteurs IoT.
Fonctionnalités du connecteur
| Élément | Valeur |
|---|---|
| Type de table | Destination |
| Mode d'exécution | Streaming |
| Format des données | Point |
| Type d'API | SQL |
| Mise à jour ou suppression des données dans la table de destination | Non pris en charge |
| Métriques | numRecordsOut, numRecordsOutPerSecond, currentSendTime |
Pour plus de détails sur ces métriques, consultez la documentation relative aux Métriques de surveillance.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Créé une base de données dans TSDB for InfluxDB. Consultez la rubrique Gérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de données
Limites
Le connecteur TSDB for InfluxDB est uniquement pris en charge par les déploiements Realtime Compute for Apache Flink utilisant VVR 2.1.5 ou une version ultérieure.
Créer une table de destination
DDL minimal
L'exemple suivant présente les colonnes minimales requises pour définir une table de destination :
CREATE TABLE influxdb_sink (
`metric` VARCHAR,
`timestamp` BIGINT,
`tag_value1` VARCHAR,
`field_fieldValue1` DOUBLE
) WITH (
'connector' = 'influxdb',
'url' = 'http://service.cn.influxdb.aliyuncs.com:****',
'database' = '<yourDatabaseName>',
'username' = '<yourDatabaseUserName>',
'password' = '<yourDatabasePassword>'
);
Conventions de nommage des colonnes du schéma
Les colonnes de la table de destination doivent respecter une convention de nommage fixe correspondant au modèle de données InfluxDB. L'ordre des colonnes est immuable.
| Position | Nom de colonne | Type | Obligatoire | Correspond à |
|---|---|---|---|---|
| 0 | metric |
VARCHAR | Oui | Nom de la mesure InfluxDB |
| 1 | timestamp |
BIGINT | Oui | Horodatage InfluxDB ; l'unité doit être la milliseconde |
| 2+ | tag_<name> |
VARCHAR | Au moins un | Tag InfluxDB (métadonnées indexées) |
| 3+ | field_<name> |
Tout type pris en charge | Au moins un | Champ InfluxDB (valeur de donnée) |
Pour écrire dans plusieurs colonnes de champ, définissez-les selon le modèle suivant :
`field_fieldValue1` DOUBLE,
`field_fieldValue2` INTEGER,
`field_fieldValueN` INTEGER
Seuls les noms de colonne metric, timestamp, tag_* et field_* sont pris en charge. Tout autre nom de colonne provoquera une erreur.
Options du connecteur
| Paramètre | Obligatoire | Valeur par défaut | Type | Description |
|---|---|---|---|---|
connector |
Oui | — | String | Doit être influxdb. |
url |
Oui | — | String | Endpoint VPC de l'instance TSDB for InfluxDB. Les protocoles HTTP et HTTPS sont pris en charge. Exemple : https://localhost:8086 ou http://localhost:3242. |
database |
Oui | — | String | Nom de la base de données. Exemple : db-flink. |
username |
Oui | — | String | Nom d'utilisateur pour la base de données. L'utilisateur doit disposer des autorisations d'écriture sur la base de données cible. Consultez la rubrique Gérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de données. |
password |
Oui | — | String | Mot de passe de l'utilisateur spécifié. Consultez la rubrique Gérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de données. |
batchSize |
Non | 300 |
Integer | Nombre d'enregistrements à écrire en un seul lot. |
retentionPolicy |
Non | autogen |
String | Politique de rétention de la base de données cible. Si ce paramètre n'est pas spécifié, la politique de rétention par défaut de la base de données (autogen) est utilisée. Consultez la rubrique Gérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de donnéesGérer les comptes utilisateur et les bases de données. |
ignoreErrorData |
Non | false |
Boolean | Gestion des erreurs d'écriture. true : ignore les erreurs d'écriture et poursuit l'exécution. false : échoue le job en cas d'erreur d'écriture. |
Mappages de types de données
| Type InfluxDB | Type Flink |
|---|---|
| BOOLEAN | BOOLEAN |
| INT | INT |
| BIGINT | BIGINT |
| FLOAT | FLOAT |
| DECIMAL | DECIMAL |
| DOUBLE | DOUBLE |
| DATE | DATE |
| TIME | TIME |
| TIMESTAMP | TIMESTAMP |
| VARCHAR | VARCHAR |
Exemple
L'exemple suivant génère des données aléatoires à l'aide du connecteur datagen et les écrit dans TSDB for InfluxDB.
CREATE TEMPORARY TABLE datagen_source (
`metric` VARCHAR,
`timestamp` BIGINT,
`fieldvalue` DOUBLE,
`tagvalue` VARCHAR
) WITH (
'connector' = 'datagen',
'fields.metric.length' = '3',
'fields.tagvalue.length' = '3',
'fields.timestamp.min' = '1587539547000',
'fields.timestamp.max' = '1619075547000',
'fields.fieldvalue.min' = '1',
'fields.fieldvalue.max' = '100000',
'rows-per-second' = '50'
);
CREATE TEMPORARY TABLE influxdb_sink (
`metric` VARCHAR,
`timestamp` BIGINT,
`field_fieldValue1` DOUBLE,
`tag_value1` VARCHAR
) WITH (
'connector' = 'influxdb',
'url' = 'https://***********.influxdata.tsdb.aliyuncs.com:****',
'database' = '<yourDatabaseName>',
'username' = '<yourDatabaseUserName>',
'password' = '<yourDatabasePassword>',
'batchSize' = '100',
'retentionPolicy' = 'autogen',
'ignoreErrorData' = 'false'
);
INSERT INTO influxdb_sink
SELECT
`metric`,
`timestamp`,
`fieldvalue`,
`tagvalue`
FROM datagen_source;