Официальный клиент ClickHouse для Rust, изначально разработанный Paul Loyd. Исходный код клиента доступен в репозитории GitHub.
Обзор
- Использует
serdeдля кодирования и декодирования строк. - Поддерживает атрибуты
serde:skip_serializing,skip_deserializing,rename. - Использует формат
RowBinaryповерх HTTP-транспорта.- Планируется переход на
Nativeповерх TCP.
- Планируется переход на
- Поддерживает TLS (через возможности
native-tlsиrustls-tls). - Поддерживает сжатие и распаковку (LZ4).
- Предоставляет API для выборки и вставки данных, выполнения DDL-запросов и батчинга на стороне клиента.
- Предоставляет удобные моки для модульного тестирования.
Установка
Чтобы использовать крейт, добавьте следующее в Cargo.toml:
[dependencies]
clickhouse = "0.12.2"
[dev-dependencies]
clickhouse = { version = "0.12.2", features = ["test-util"] }См. также: страница на crates.io.
Возможности Cargo
lz4(включена по умолчанию) — включает вариантыCompression::Lz4иCompression::Lz4Hc(_). Если эта возможность включена,Compression::Lz4по умолчанию используется для всех запросов, кромеWATCH.native-tls— поддерживает URL со схемойHTTPSчерезhyper-tls, который собирается с OpenSSL.rustls-tls— поддерживает URL со схемойHTTPSчерезhyper-rustls, который не собирается с OpenSSL.inserter— включаетclient.inserter().test-util— добавляет моки. См. пример. Используйте только вdev-dependencies.watch— включает функциональностьclient.watch. Подробности см. в соответствующем разделе.uuid— добавляетserde::uuidдля работы с крейтом uuid.time— добавляетserde::timeдля работы с крейтом time.
Совместимость версий ClickHouse
Клиент совместим с LTS-версиями ClickHouse и более новыми версиями, а также с ClickHouse Cloud.
ClickHouse server версий ниже v22.6 в некоторых редких случаях некорректно обрабатывает RowBinary.
Чтобы решить эту проблему, можно использовать v0.11+ и включить возможность wa-37420. Примечание: эту возможность не следует использовать с более новыми версиями ClickHouse.
Примеры
Мы стараемся охватить различные сценарии использования клиента в примерах в репозитории клиента. Обзор доступен в README с примерами.
Если в примерах или в приведённой ниже документации что-то неясно или чего-то не хватает, свяжитесь с нами.
Использование
Создание экземпляра клиента
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");Подключение по HTTPS или к ClickHouse Cloud
HTTPS работает с возможностями Cargo rustls-tls и native-tls.
Затем создайте клиент как обычно. В этом примере переменные окружения используются для хранения сведений о подключении:
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"));См. также:
- Пример HTTPS с ClickHouse Cloud в репозитории client. Это также применимо к HTTPS-подключениям в собственной инфраструктуре.
Выборка строк
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? { .. }- Плейсхолдер
?fieldsзаменяется наno, name(поляRow). - Плейсхолдер
?заменяется значениями в последующих вызовахbind(). - Для получения первой строки или всех строк соответственно можно использовать удобные методы
fetch_one::<Row>()иfetch_all::<Row>(). sql::Identifierможно использовать для подстановки имён таблиц.
Примечание: поскольку весь ответ передаётся в потоковом режиме, курсоры могут возвращать ошибку даже после выдачи некоторых строк. Если в вашем случае это происходит, попробуйте query(...).with_option("wait_end_of_query", "1"), чтобы включить буферизацию ответа на стороне сервера. Подробнее. Также может быть полезен параметр buffer_size.
Вставка строк
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?;- Если
end()не вызвать,INSERTбудет прерван. - Строки отправляются постепенно, в виде потока, чтобы распределить нагрузку на сеть.
- ClickHouse выполняет атомарную вставку батчей, только если все строки помещаются в одну и ту же партицию и их количество меньше
max_insert_block_size.
Асинхронная вставка (батчинг на стороне сервера)
Вы можете использовать асинхронные вставки ClickHouse, чтобы избежать батчинга входящих данных на стороне клиента. Для этого достаточно просто передать параметр async_insert методу insert (или даже самому экземпляру Client, чтобы это применялось ко всем вызовам insert).
let client = Client::default()
.with_url("http://localhost:8123")
.with_option("async_insert", "1")
.with_option("wait_for_async_insert", "0");См. также:
- Пример async insert в репозитории клиента.
Возможность Inserter (батчинг на стороне клиента)
Требуется возможность Cargo inserter.
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?;Inserterзавершает активную вставку вcommit(), если достигнут любой из порогов (max_bytes,max_rows,period).- Интервал между завершениями активных
INSERTможно сместить с помощьюwith_period_bias, чтобы избежать пиков нагрузки при параллельной работеInserter. Inserter::time_left()можно использовать, чтобы определить, когда закончится текущий период. Если ваш поток редко выдает элементы, снова вызовитеInserter::commit(), чтобы проверить ограничения.- Пороги по времени реализованы с использованием крейта quanta, чтобы ускорить работу
inserter. Не используется, если включенtest-util(поэтому в пользовательских тестах временем можно управлять черезtokio::time::advance()). - Все строки между вызовами
commit()вставляются одним операторомINSERT.
Выполнение DDL-запросов
Для одноузлового развертывания достаточно выполнять DDL-запросы так:
client.query("DROP TABLE IF EXISTS some").execute().await?;Однако в кластерных развертываниях с балансировщиком нагрузки или в ClickHouse Cloud рекомендуется дождаться, пока DDL применится на всех репликах, используя параметр wait_end_of_query. Это можно сделать так:
client
.query("DROP TABLE IF EXISTS some")
.with_option("wait_end_of_query", "1")
.execute()
.await?;Настройки ClickHouse
Вы можете применять различные настройки ClickHouse с помощью метода with_option. Например:
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?;Помимо query, это аналогично работает и с методами insert и inserter; кроме того, тот же метод можно вызвать у экземпляра Client, чтобы задать глобальные настройки для всех запросов.
Query id
С помощью .with_option можно задать параметр query_id для идентификации запросов в журнале запросов ClickHouse.
let numbers = client
.query("SELECT number FROM system.numbers LIMIT 1")
.with_option("query_id", "some-query-id")
.fetch_all::<u64>()
.await?;Помимо query, это работает аналогичным образом и для методов insert и inserter.
См. также: пример query_id в репозитории клиента.
Идентификатор сеанса
Как и в случае с query_id, можно задать session_id, чтобы выполнять команды в рамках одного сеанса. session_id можно задать либо глобально на уровне клиента, либо для отдельных вызовов query, insert или inserter.
let client = Client::default()
.with_url("http://localhost:8123")
.with_option("session_id", "my-session");См. также: пример session_id в репозитории клиента.
Пользовательские HTTP-заголовки
Если вы используете аутентификацию через прокси или вам нужно передать пользовательские заголовки, это можно сделать так:
let client = Client::default()
.with_url("http://localhost:8123")
.with_header("X-My-Header", "hello");См. также: пример использования пользовательских HTTP-заголовков в репозитории client.
Собственный HTTP-клиент
Это может быть полезно для настройки параметров используемого пула HTTP-соединений.
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");См. также: пример пользовательского HTTP-клиента в репозитории клиента.
Типы данных
(U)Int(8|16|32|64|128)сопоставляется с соответствующими типами(u|i)(8|16|32|64|128)или типами newtype-обёртка на их основе и обратно.(U)Int256напрямую не поддерживаются, но для них есть обходное решение.Float(32|64)сопоставляется с соответствующимиf(32|64)или типами newtype-обёртка на их основе и обратно.Decimal(32|64|128)сопоставляется с соответствующимиi(32|64|128)или типами newtype-обёртка на их основе и обратно. Удобнее использоватьfixnumили другую реализацию знаковых чисел с фиксированной запятой.Booleanсопоставляется сboolили типами newtype-обёртка на его основе и обратно.Stringсопоставляется с любыми строковыми или байтовыми типами и обратно, например&str,&[u8],String,Vec<u8>илиSmartString. Пользовательские типы тоже поддерживаются. Для хранения байтов рекомендуется использоватьserde_bytes, поскольку это эффективнее.
#[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)поддерживается в виде массива байтов, например[u8; N].
#[derive(Row, Debug, Serialize, Deserialize)]
struct MyRow {
fixed_str: [u8; 16], // FixedString(16)
}Enum(8|16)поддерживается с помощьюserde_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,
}UUIDпреобразуется в/изuuid::Uuidс помощьюserde::uuid. Для этого требуется возможностьuuid.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::uuid")]
uuid: uuid::Uuid,
}IPv6преобразуется в/изstd::net::Ipv6Addr.IPv4преобразуется в/изstd::net::Ipv4Addrс помощьюserde::ipv4.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::ipv4")]
ipv4: std::net::Ipv4Addr,
}Dateпреобразуется в/изu16илиnewtype-обёрткуна его основе и представляет собой количество дней, прошедших с1970-01-01. Также поддерживаетсяtime::Dateчерезserde::time::date, для чего требуется возможностьtime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
days: u16,
#[serde(with = "clickhouse::serde::time::date")]
date: Date,
}Date32преобразуется в/изi32илиnewtype-обёрткуна его основе и представляет количество дней, прошедших с1970-01-01. Также поддерживаетсяtime::Dateчерезserde::time::date32; для этого требуется возможностьtime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
days: i32,
#[serde(with = "clickhouse::serde::time::date32")]
date: Date,
}DateTimeсопоставляется сu32и обратно, либо с newtype-обёрткой над ним, и представляет собой количество секунд, прошедших с эпохи Unix. Также поддерживаетсяtime::OffsetDateTimeчерезserde::time::datetime, для чего требуется возможностьtime.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
ts: u32,
#[serde(with = "clickhouse::serde::time::datetime")]
dt: OffsetDateTime,
}DateTime64(_)преобразуется в/изi32илиnewtype, оборачивающий его, и представляет время, прошедшее с эпохи Unix. Также поддерживаетсяtime::OffsetDateTimeпри использованииserde::time::datetime64::*; для этого требуется возможность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, ...)преобразуется в(A, B, ...)и обратно, либо в newtype-обёртку над ним.Array(_)преобразуется в любой срез и обратно, напримерVec<_>,&[_]. Также поддерживаются новые типы.Map(K, V)ведёт себя какArray((K, V)).LowCardinality(_)поддерживается прозрачно.Nullable(_)преобразуется вOption<_>и обратно. Для хелперовclickhouse::serde::*добавьте::option.
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::ipv4::option")]
ipv4_opt: Option<Ipv4Addr>,
}Nestedподдерживается через передачу нескольких массивов с переименованием.
// 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>,
}- Поддерживаются типы
Geo.Pointведёт себя как кортеж(f64, f64), а остальные типы — это просто срезы точек.
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,
}- Типы данных
Variant,Dynamicи (новый)JSONпока не поддерживаются.
Мокирование
Крейт предоставляет утилиты для мокирования сервера CH и тестирования запросов DDL, SELECT, INSERT и WATCH. Эту функциональность можно включить, активировав возможность test-util. Используйте её только как зависимость для разработки.
См. пример.
Устранение неполадок
CANNOT_READ_ALL_DATA
Наиболее частая причина ошибки CANNOT_READ_ALL_DATA заключается в том, что определение строки на стороне приложения не совпадает с определением в ClickHouse.
Рассмотрим следующую таблицу:
CREATE OR REPLACE TABLE event_log (id UInt32)
ENGINE = MergeTree
ORDER BY timestampЗатем, если EventLog определён на стороне приложения с несовпадающими типами, например:
#[derive(Debug, Serialize, Deserialize, Row)]
struct EventLog {
id: String, // <- should be u32 instead!
}При вставке данных может возникнуть следующая ошибка:
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)")В этом примере проблема исправляется правильным определением структуры EventLog:
#[derive(Debug, Serialize, Deserialize, Row)]
struct EventLog {
id: u32
}Известные ограничения
- Типы данных
Variant,Dynamic, (new)JSONпока не поддерживаются. - Привязка параметров на стороне сервера пока не поддерживается; отслеживать статус можно в этой задаче.
Свяжитесь с нами
Если у вас есть вопросы или вам нужна помощь, обращайтесь к нам в Community Slack или через GitHub Issues.