Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Клиент ClickHouse для Rust

Официальный клиент 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");

См. также:

Возможность 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,
}
#[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.

Navigation