Le client Rust officiel pour se connecter à ClickHouse, développé à l’origine par Paul Loyd. Le code source du client est disponible dans le dépôt GitHub.
Vue d’ensemble
- Utilise
serdepour l’encodage/décodage des lignes. - Prend en charge les attributs
serde:skip_serializing,skip_deserializing,rename. - Utilise le format
RowBinaryvia HTTP.- Un passage à
Nativesur TCP est prévu.
- Un passage à
- Prend en charge TLS (via les fonctionnalités
native-tlsetrustls-tls). - Prend en charge la compression et la décompression (LZ4).
- Fournit des API pour interroger ou insérer des données, exécuter des DDL et effectuer du batching côté client.
- Fournit des mocks pratiques pour les tests unitaires.
Installation
Pour utiliser ce crate, ajoutez ce qui suit à votre fichier Cargo.toml :
[dependencies]
clickhouse = "0.12.2"
[dev-dependencies]
clickhouse = { version = "0.12.2", features = ["test-util"] }Voir aussi : page sur crates.io.
Fonctionnalités Cargo
lz4(activée par défaut) — active les variantesCompression::Lz4etCompression::Lz4Hc(_). Si elle est activée,Compression::Lz4est utilisée par défaut pour toutes les requêtes, saufWATCH.native-tls— prend en charge les URL utilisant le schémaHTTPSviahyper-tls, qui s'appuie sur OpenSSL.rustls-tls— prend en charge les URL utilisant le schémaHTTPSviahyper-rustls, qui ne s'appuie pas sur OpenSSL.inserter— activeclient.inserter().test-util— ajoute des mocks. Voir l'exemple. À utiliser uniquement dansdev-dependencies.watch— active la fonctionnalitéclient.watch. Voir la section correspondante pour plus de détails.uuid— ajouteserde::uuidpour utiliser la crate uuid.time— ajouteserde::timepour utiliser la crate time.
Compatibilité des versions de ClickHouse
Le client est compatible avec les versions LTS de ClickHouse, ainsi qu'avec les versions plus récentes, et avec ClickHouse Cloud.
Le serveur ClickHouse antérieur à la version v22.6 gère RowBinary incorrectement dans certains rares cas.
Vous pouvez utiliser v0.11+ et activer la fonctionnalité wa-37420 pour corriger ce problème. Remarque : cette fonctionnalité ne doit pas être utilisée avec des versions plus récentes de ClickHouse.
Exemples
Nous cherchons à couvrir différents scénarios d'utilisation du client à l'aide des examples du dépôt du client. La vue d'ensemble est disponible dans le README des examples.
Si certains points ne sont pas clairs ou s'il manque des éléments dans les examples ou dans la documentation ci-dessous, n'hésitez pas à nous contacter.
Utilisation
Créer une instance de client
use clickhouse::Client;
let client = Client::default()
// should include both protocol and port
.with_url("http://localhost:8123")
.with_user("name")
.with_password("123")
.with_database("test");Connexion HTTPS ou à ClickHouse Cloud
HTTPS fonctionne avec les fonctionnalités Cargo rustls-tls ou native-tls.
Créez ensuite le client comme d’habitude. Dans cet exemple, les variables d’environnement servent à stocker les informations de connexion :
fn read_env_var(key: &str) -> String {
env::var(key).unwrap_or_else(|_| panic!("{key} env variable should be set"))
}
let client = Client::default()
.with_url(read_env_var("CLICKHOUSE_URL"))
.with_user(read_env_var("CLICKHOUSE_USER"))
.with_password(read_env_var("CLICKHOUSE_PASSWORD"));Voir aussi :
- Exemple HTTPS avec ClickHouse Cloud dans le dépôt du client. Cela devrait également s’appliquer aux connexions HTTPS on-premise.
Sélectionner des lignes
use serde::Deserialize;
use clickhouse::Row;
use clickhouse::sql::Identifier;
#[derive(Row, Deserialize)]
struct MyRow<'a> {
no: u32,
name: &'a str,
}
let table_name = "some";
let mut cursor = client
.query("SELECT ?fields FROM ? WHERE no BETWEEN ? AND ?")
.bind(Identifier(table_name))
.bind(500)
.bind(504)
.fetch::<MyRow<'_>>()?;
while let Some(row) = cursor.next().await? { .. }- L’espace réservé
?fieldsest remplacé parno, name(champs deRow). - L’espace réservé
?est remplacé par les valeurs des appelsbind()suivants. - Les méthodes pratiques
fetch_one::<Row>()etfetch_all::<Row>()peuvent être utilisées pour récupérer respectivement la première ligne ou toutes les lignes. sql::Identifierpeut être utilisé pour lier des noms de table.
NB : comme la réponse entière est transmise en flux, les curseurs peuvent renvoyer une erreur même après avoir produit certaines lignes. Si cela se produit dans votre cas d’usage, vous pouvez essayer query(...).with_option("wait_end_of_query", "1") afin d’activer la mise en mémoire tampon de la réponse côté serveur. Plus de détails. L’option buffer_size peut également être utile.
Insertion de lignes
use serde::Serialize;
use clickhouse::Row;
#[derive(Row, Serialize)]
struct MyRow {
no: u32,
name: String,
}
let mut insert = client.insert("some")?;
insert.write(&MyRow { no: 0, name: "foo".into() }).await?;
insert.write(&MyRow { no: 1, name: "bar".into() }).await?;
insert.end().await?;- Si
end()n’est pas appelé, l’INSERTest annulé. - Les lignes sont envoyées progressivement en flux afin de répartir la charge sur le réseau.
- ClickHouse n’insère les lots de façon atomique que si toutes les lignes tiennent dans la même partition et que leur nombre est inférieur à
max_insert_block_size.
Async insert (batching côté serveur)
Vous pouvez utiliser les insertions asynchrones de ClickHouse pour éviter le batching côté client des données entrantes. Pour cela, il suffit de spécifier l'option async_insert pour la méthode insert (ou même pour l'instance Client elle-même, afin qu'elle s'applique à tous les appels à insert).
let client = Client::default()
.with_url("http://localhost:8123")
.with_option("async_insert", "1")
.with_option("wait_for_async_insert", "0");Voir aussi :
- Exemple d’utilisation d’async insert dans le dépôt client.
Fonctionnalité Inserter (batching côté client)
Nécessite la fonctionnalité inserter de Cargo.
let mut inserter = client.inserter("some")?
.with_timeouts(Some(Duration::from_secs(5)), Some(Duration::from_secs(20)))
.with_max_bytes(50_000_000)
.with_max_rows(750_000)
.with_period(Some(Duration::from_secs(15)));
inserter.write(&MyRow { no: 0, name: "foo".into() })?;
inserter.write(&MyRow { no: 1, name: "bar".into() })?;
let stats = inserter.commit().await?;
if stats.rows > 0 {
println!(
"{} bytes, {} rows, {} transactions have been inserted",
stats.bytes, stats.rows, stats.transactions,
);
}
// don't forget to finalize the inserter during the application shutdown
// and commit the remaining rows. `.end()` will provide stats as well.
inserter.end().await?;Insertertermine l’insertion active danscommit()si l’un des seuils (max_bytes,max_rows,period) est atteint.- L’intervalle entre la fin des
INSERTactifs peut être ajusté à l’aide dewith_period_biasafin d’éviter des pics de charge causés par des inserters parallèles. Inserter::time_left()peut être utilisé pour détecter quand la période en cours se termine. Appelez à nouveauInserter::commit()pour vérifier les limites si votre flux émet rarement des éléments.- Les seuils de temps sont implémentés à l’aide de la crate quanta pour accélérer
inserter. Ils ne sont pas utilisés sitest-utilest activé (le temps peut alors être géré viatokio::time::advance()dans des tests personnalisés). - Toutes les lignes entre deux appels à
commit()sont insérées dans la même instructionINSERT.
Exécution des DDL
Avec un déploiement sur un seul nœud, il suffit d’exécuter les DDL comme suit :
client.query("DROP TABLE IF EXISTS some").execute().await?;Cependant, dans les déploiements en cluster avec un équilibreur de charge ou ClickHouse Cloud, il est recommandé d'attendre que le DDL soit appliqué à toutes les répliques à l'aide de l'option wait_end_of_query. Voici comment procéder :
client
.query("DROP TABLE IF EXISTS some")
.with_option("wait_end_of_query", "1")
.execute()
.await?;Paramètres ClickHouse
Vous pouvez appliquer différents paramètres ClickHouse à l’aide de la méthode with_option. Par exemple :
let numbers = client
.query("SELECT number FROM system.numbers")
// This setting will be applied to this particular query only;
// it will override the global client setting.
.with_option("limit", "3")
.fetch_all::<u64>()
.await?;En plus de query, cela fonctionne de la même manière avec les méthodes insert et inserter ; on peut également appeler cette même méthode sur l’instance Client afin de définir des paramètres globaux pour toutes les requêtes.
ID de requête
Avec .with_option, vous pouvez définir l’option query_id pour identifier les requêtes dans le journal des requêtes de ClickHouse.
let numbers = client
.query("SELECT number FROM system.numbers LIMIT 1")
.with_option("query_id", "some-query-id")
.fetch_all::<u64>()
.await?;Outre query, cela fonctionne de façon similaire avec les méthodes insert et inserter.
Voir aussi : l’exemple query_id dans le dépôt du client.
ID de session
Comme pour query_id, vous pouvez définir session_id afin d’exécuter les instructions dans la même session. session_id peut être défini soit globalement au niveau du client, soit pour chaque appel à query, insert ou inserter.
let client = Client::default()
.with_url("http://localhost:8123")
.with_option("session_id", "my-session");Voir aussi : l’exemple session_id dans le dépôt du client.
En-têtes HTTP personnalisés
Si vous utilisez l’authentification via un proxy ou si vous devez transmettre des en-têtes personnalisés, vous pouvez procéder comme suit :
let client = Client::default()
.with_url("http://localhost:8123")
.with_header("X-My-Header", "hello");Voir aussi : exemple d’en-têtes HTTP personnalisés dans le dépôt du client.
Client HTTP personnalisé
Cela peut être utile pour ajuster les paramètres du pool de connexions HTTP sous-jacent.
use hyper_util::client::legacy::connect::HttpConnector;
use hyper_util::client::legacy::Client as HyperClient;
use hyper_util::rt::TokioExecutor;
let connector = HttpConnector::new(); // or HttpsConnectorBuilder
let hyper_client = HyperClient::builder(TokioExecutor::new())
// For how long keep a particular idle socket alive on the client side (in milliseconds).
// It is supposed to be a fair bit less that the ClickHouse server KeepAlive timeout,
// which was by default 3 seconds for pre-23.11 versions, and 10 seconds after that.
.pool_idle_timeout(Duration::from_millis(2_500))
// Sets the maximum idle Keep-Alive connections allowed in the pool.
.pool_max_idle_per_host(4)
.build(connector);
let client = Client::with_http_client(hyper_client).with_url("http://localhost:8123");Voir aussi : exemple de client HTTP personnalisé dans le dépôt du client.
Types de données
(U)Int(8|16|32|64|128)correspond aux types(u|i)(8|16|32|64|128)correspondants, ou à desnewtypesqui les encapsulent, et inversement.(U)Int256ne sont pas pris en charge nativement, mais il existe une solution de contournement.Float(32|64)correspond aux typesf(32|64)correspondants, ou à desnewtypesqui les encapsulent, et inversement.Decimal(32|64|128)correspond aux typesi(32|64|128)correspondants, ou à desnewtypesqui les encapsulent, et inversement. Il est plus pratique d'utiliserfixnumou une autre implémentation de nombres à virgule fixe signés.Booleancorrespond àboolou à desnewtypesqui l'encapsulent, et inversement.Stringcorrespond à n'importe quel type de chaîne ou d'octets, par exemple&str,&[u8],String,Vec<u8>ouSmartString. Lesnewtypessont également pris en charge. Pour stocker des octets, envisagez d'utiliserserde_bytes, car c'est plus efficace.
#[derive(Row, Debug, Serialize, Deserialize)]
struct MyRow<'a> {
str: &'a str,
string: String,
#[serde(with = "serde_bytes")]
bytes: Vec<u8>,
#[serde(with = "serde_bytes")]
byte_slice: &'a [u8],
}FixedString(N)est pris en charge sous la forme d’un tableau d’octets, par ex.[u8; N].
#[derive(Row, Debug, Serialize, Deserialize)]
struct MyRow {
fixed_str: [u8; 16], // FixedString(16)
}Enum(8|16)est pris en charge viaserde_repr.
use serde_repr::{Deserialize_repr, Serialize_repr};
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
level: Level,
}
#[derive(Debug, Serialize_repr, Deserialize_repr)]
#[repr(u8)]
enum Level {
Debug = 1,
Info = 2,
Warn = 3,
Error = 4,
}UUIDse convertit depuis/versuuid::Uuidà l’aide deserde::uuid. Nécessite la fonctionnalitéuuid.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::uuid")]
uuid: uuid::Uuid,
}IPv6est associé àstd::net::Ipv6Addr, et inversement.IPv4est associé àstd::net::Ipv4Addr, et inversement, viaserde::ipv4.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::ipv4")]
ipv4: std::net::Ipv4Addr,
}Datese convertit depuis/versu16ou un newtype qui l’encapsule, et représente un nombre de jours écoulés depuis1970-01-01. De plus,time::Dateest également pris en charge viaserde::time::date, ce qui requiert la fonctionnalitétime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
days: u16,
#[serde(with = "clickhouse::serde::time::date")]
date: Date,
}Date32se mappe depuis/versi32ou un newtype qui l'encapsule, et représente un nombre de jours écoulés depuis1970-01-01. De plus,time::Dateest également pris en charge viaserde::time::date32, ce qui nécessite la featuretime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
days: i32,
#[serde(with = "clickhouse::serde::time::date32")]
date: Date,
}DateTimese convertit depuis/versu32ou un newtype qui l’encapsule, et représente un nombre de secondes écoulées depuis l’époque Unix. De plus,time::OffsetDateTimeest également pris en charge viaserde::time::datetime, ce qui nécessite la fonctionnalitétime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
ts: u32,
#[serde(with = "clickhouse::serde::time::datetime")]
dt: OffsetDateTime,
}DateTime64(_)se convertit vers/depuisi32ou un newtype qui l’encapsule, et représente un temps écoulé depuis l’époque Unix. De plus,time::OffsetDateTimeest également pris en charge viaserde::time::datetime64::*, ce qui nécessite la fonctionnalitétime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
ts: i64, // elapsed s/us/ms/ns depending on `DateTime64(X)`
#[serde(with = "clickhouse::serde::time::datetime64::secs")]
dt64s: OffsetDateTime, // `DateTime64(0)`
#[serde(with = "clickhouse::serde::time::datetime64::millis")]
dt64ms: OffsetDateTime, // `DateTime64(3)`
#[serde(with = "clickhouse::serde::time::datetime64::micros")]
dt64us: OffsetDateTime, // `DateTime64(6)`
#[serde(with = "clickhouse::serde::time::datetime64::nanos")]
dt64ns: OffsetDateTime, // `DateTime64(9)`
}Tuple(A, B, ...)correspond à(A, B, ...)et vice versa, ou à un newtype qui l’encapsule.Array(_)correspond à n’importe quelle slice et vice versa, par ex.Vec<_>,&[_]. Les nouveaux types sont également pris en charge.Map(K, V)se comporte commeArray((K, V)).LowCardinality(_)est pris en charge de manière transparente.Nullable(_)correspond àOption<_>et vice versa. Pour les helpersclickhouse::serde::*, ajoutez::option.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::ipv4::option")]
ipv4_opt: Option<Ipv4Addr>,
}Nestedest pris en charge en fournissant plusieurs tableaux renommés.
// CREATE TABLE test(items Nested(name String, count UInt32))
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(rename = "items.name")]
items_name: Vec<String>,
#[serde(rename = "items.count")]
items_count: Vec<u32>,
}- Les types
Geosont pris en charge.Pointse comporte comme un uplet(f64, f64), et les autres types ne sont que des slices de points.
type Point = (f64, f64);
type Ring = Vec<Point>;
type Polygon = Vec<Ring>;
type MultiPolygon = Vec<Polygon>;
type LineString = Vec<Point>;
type MultiLineString = Vec<LineString>;
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
point: Point,
ring: Ring,
polygon: Polygon,
multi_polygon: MultiPolygon,
line_string: LineString,
multi_line_string: MultiLineString,
}- Les types de données
Variant,DynamicetJSON(nouveau) ne sont pas encore pris en charge.
Simulation
Le crate fournit des utilitaires permettant de simuler un serveur CH et de tester les requêtes DDL, SELECT, INSERT et WATCH. Cette fonctionnalité peut être activée avec la fonctionnalité test-util. Utilisez-la uniquement comme dépendance de développement.
Voir l’exemple.
Dépannage
CANNOT_READ_ALL_DATA
La cause la plus fréquente de l’erreur CANNOT_READ_ALL_DATA est que la définition de la ligne côté application ne correspond pas à celle de ClickHouse.
Considérez la table suivante :
CREATE OR REPLACE TABLE event_log (id UInt32)
ENGINE = MergeTree
ORDER BY timestampEnsuite, si EventLog est défini du côté de l’application avec des types incompatibles, par exemple :
#[derive(Debug, Serialize, Deserialize, Row)]
struct EventLog {
id: String, // <- should be u32 instead!
}Lors de l’insertion des données, l’erreur suivante peut survenir :
Error: BadResponse("Code: 33. DB::Exception: Cannot read all data. Bytes read: 5. Bytes expected: 23.: (at row 1)\n: While executing BinaryRowInputFormat. (CANNOT_READ_ALL_DATA)")Dans cet exemple, ce problème est corrigé grâce à la définition correcte de la struct EventLog :
#[derive(Debug, Serialize, Deserialize, Row)]
struct EventLog {
id: u32
}Limites connues
- Les types de données
Variant,DynamicetJSON(nouveau) ne sont pas encore pris en charge. - La liaison des paramètres côté serveur n'est pas encore prise en charge ; voir ce ticket pour en suivre l'avancement.
Contactez-nous
Si vous avez des questions ou avez besoin d’aide, n’hésitez pas à nous contacter sur le Slack de la communauté ou via les issues GitHub.