Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

ClickHouse Rust クライアント

ClickHouse に接続するための公式 Rust クライアントです。もともとは Paul Loyd によって開発されました。クライアントのソースコードは GitHub リポジトリ で公開されています。

概要

  • 行のエンコード/デコードに serde を使用します。
  • serde の属性 skip_serializingskip_deserializingrename をサポートしています。
  • HTTP 経由で RowBinary フォーマットを使用します。
    • 将来的には、TCP 上の Native への切り替えが予定されています。
  • TLS (native-tls および rustls-tls の feature 経由) をサポートしています。
  • 圧縮および解凍 (LZ4) をサポートしています。
  • データの取得や挿入、DDL の実行、クライアント側でのバッチ処理のための API を提供します。
  • 単体テストに便利なモックを提供します。

インストール

この クレート を使用するには、次の内容を 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(_) バリアントを有効にします。有効にすると、WATCH を除くすべてのクエリで既定で Compression::Lz4 が使用されます。
  • native-tls — OpenSSL にリンクする hyper-tls を介して、HTTPS スキームの URL をサポートします。
  • rustls-tls — OpenSSL にリンクしない hyper-rustls を介して、HTTPS スキームの URL をサポートします。
  • inserterclient.inserter() を有効にします。
  • test-util — モックを追加します。 を参照してください。使用は dev-dependencies のみにしてください。
  • watchclient.watch 機能を有効にします。詳細は該当する節を参照してください。
  • uuiduuid クレートと連携するために serde::uuid を追加します。
  • timetime クレートと連携するために serde::time を追加します。

ClickHouse バージョンの互換性

このクライアントは、ClickHouse の LTS 版以降のバージョンおよび ClickHouse Cloud と互換性があります。

v22.6 より前の ClickHouse server では、まれに RowBinary が正しく処理されない場合があります。 この問題を回避するには、v0.11 以降を使用し、wa-37420 機能を有効にしてください。注: この機能は、より新しい ClickHouse バージョンでは使用しないでください。

クライアントの利用に関するさまざまなシナリオを、クライアントリポジトリ内の examples でカバーすることを目指しています。概要については、examples README を参照してください。

examples や以下のドキュメントで不明な点や不足している内容があれば、お気軽にお問い合わせください

使い方

クライアントインスタンスの作成

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 は、rustls-tls または native-tls のいずれの Cargo feature でも利用できます。

次に、通常どおり client を作成します。この例では、接続情報を保存するために環境変数を使用します。

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"));

関連項目:

行を選択する

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? { .. }
  • プレースホルダー ?fieldsno, 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 未満の場合に限られます。

非同期 INSERT (サーバー側バッチ処理)

受信データをクライアント側でバッチ処理しないようにするには、ClickHouse の非同期挿入を使用できます。これは、insert メソッドに async_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 機能 (クライアント側バッチ処理)

inserter Cargo feature が必要です。

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?;
  • いずれかのしきい値 (max_bytesmax_rowsperiod) に達すると、Insertercommit() 内でアクティブな insert を終了します。
  • アクティブな INSERT の終了間隔は、並列 inserter による負荷のスパイクを避けるために、with_period_bias を使って偏らせることができます。
  • Inserter::time_left() は、現在の period がいつ終了するかを検出するために使用できます。ストリームがまれにしか項目を送出しない場合は、Inserter::commit() を再度呼び出して制限値を確認してください。
  • 時間しきい値は、inserter を高速化するために quanta クレート を使って実装されています。test-util が enabled の場合は使用されません (そのため、カスタムテストでは tokio::time::advance() で time を管理できます) 。
  • commit() 呼び出しの間にあるすべての行は、同じ INSERT ステートメントに挿入されます。

DDL の実行

単一ノード構成では、次のように DDL を実行すれば十分です。

client.query("DROP TABLE IF EXISTS some").execute().await?;

ただし、ロードバランサーを使用するクラスター構成のデプロイメントや ClickHouse Cloud では、wait_end_of_query オプションを使って、すべてのレプリカに DDL が適用されるのを待つことを推奨します。次のように実行できます。

client
    .query("DROP TABLE IF EXISTS some")
    .with_option("wait_end_of_query", "1")
    .execute()
    .await?;

ClickHouseの設定

with_option メソッドを使うと、さまざまな ClickHouseの設定 を適用できます。例:

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インスタンスに対して呼び出すこともできます。

クエリ 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 メソッドでも機能します。

関連項目: client リポジトリの query_id の例

セッションID

query_id と同様に、session_id を設定すると、同じセッションでステートメントを実行できます。session_id は、クライアントレベルでグローバルに設定することも、queryinsert、または inserter の呼び出しごとに設定することもできます。

let client = Client::default()
    .with_url("http://localhost:8123")
    .with_option("session_id", "my-session");

関連項目: client リポジトリの session_id の例

カスタムHTTPヘッダー

プロキシ認証を使用している場合や、カスタムHTTPヘッダーを渡す必要がある場合は、次のように設定できます。

let client = Client::default()
    .with_url("http://localhost:8123")
    .with_header("X-My-Header", "hello");

関連項目: client リポジトリにあるカスタム HTTP ヘッダーの例

カスタム 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");

関連項目: clientリポジトリのcustom HTTP client example

データ型

  • (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 や、その他の符号付き固定小数点数実装を使うほうが便利です。
  • Booleanbool、またはそれをラップした newtype と相互変換できます。
  • String は任意の文字列型またはバイト列型と相互変換できます。たとえば &str&[u8]StringVec<u8>SmartString です。newtype もサポートされています。バイト列を保存する場合は、より効率的な 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,
}
  • UUIDserde::uuid を使用することで uuid::Uuid との間で相互変換されます。uuid feature が必要です。
#[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,
}
  • Dateu16 またはそれをラップした newtype と相互変換でき、1970-01-01 からの経過日数を表します。また、time::Dateserde::time::date を使用することでサポートされますが、その場合は time フィーチャーが必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
    days: u16,
    #[serde(with = "clickhouse::serde::time::date")]
    date: Date,
}
  • Date32i32 またはそれをラップした newtype に対応し、1970-01-01 からの経過日数を表します。また、time::Dateserde::time::date32 を使うことでサポートされますが、これには time feature が必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
    days: i32,
    #[serde(with = "clickhouse::serde::time::date32")]
    date: Date,
}
  • DateTimeu32 またはそれをラップした newtype と相互変換でき、UNIX epoch からの経過秒数を表します。また、time::OffsetDateTimeserde::time::datetime を使うことでサポートされますが、これには time feature が必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
    ts: u32,
    #[serde(with = "clickhouse::serde::time::datetime")]
    dt: OffsetDateTime,
}
  • DateTime64(_)i32 またはそれをラップする newtype に相互変換され、Unix epoch からの経過時間を表します。また、time::OffsetDateTimeserde::time::datetime64::* を使用することでサポートされますが、これには time feature が必要です。
#[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) というタプルとして扱われ、それ以外の型は単に Point のスライスです。
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,
}
  • VariantDynamic、 (新しい) JSON データ型はまだサポートされていません。

モック化

このクレートには、CHサーバーをモック化し、DDL、SELECTINSERTWATCH クエリをテストするためのユーティリティが用意されています。この機能は test-util feature で有効にできます。必ず開発用依存関係としてのみ使用してください。

サンプル を参照してください。

トラブルシューティング

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 struct を正しく定義することで、この問題を修正できます。

#[derive(Debug, Serialize, Deserialize, Row)]
struct EventLog {
    id: u32
}

既知の制限事項

  • VariantDynamic、 (新しい) JSON データ型は、まだサポートされていません。
  • サーバー側のパラメータバインドは、まだサポートされていません。追跡状況については、この issue を参照してください。

お問い合わせ

ご不明な点やサポートが必要な場合は、Community Slack または GitHub issues でお気軽にご連絡ください。

Navigation