ClickHouse に接続するための公式 Rust クライアントです。もともとは Paul Loyd によって開発されました。クライアントのソースコードは GitHub リポジトリ で公開されています。
概要
- 行のエンコード/デコードに
serdeを使用します。 serdeの属性skip_serializing、skip_deserializing、renameをサポートしています。- HTTP 経由で
RowBinaryフォーマットを使用します。- 将来的には、TCP 上の
Nativeへの切り替えが予定されています。
- 将来的には、TCP 上の
- 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 をサポートします。inserter—client.inserter()を有効にします。test-util— モックを追加します。例 を参照してください。使用はdev-dependenciesのみにしてください。watch—client.watch機能を有効にします。詳細は該当する節を参照してください。uuid— uuid クレートと連携するためにserde::uuidを追加します。time— time クレートと連携するために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"));関連項目:
- クライアントリポジトリのClickHouse Cloud での HTTPS の例。これはオンプレミス環境での 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未満の場合に限られます。
非同期 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");関連項目:
- 非同期 INSERT の例 (client リポジトリ) 。
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_bytes、max_rows、period) に達すると、Inserterはcommit()内でアクティブな 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 は、クライアントレベルでグローバルに設定することも、query、insert、または 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や、その他の符号付き固定小数点数実装を使うほうが便利です。Booleanはbool、またはそれをラップした newtype と相互変換できます。Stringは任意の文字列型またはバイト列型と相互変換できます。たとえば&str、&[u8]、String、Vec<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,
}UUIDはserde::uuidを使用することでuuid::Uuidとの間で相互変換されます。uuidfeature が必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
#[serde(with = "clickhouse::serde::uuid")]
uuid: uuid::Uuid,
}IPv6はstd::net::Ipv6Addrと相互に変換されます。IPv4はserde::ipv4を使用してstd::net::Ipv4Addrと相互に変換されます。
#[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を使うことでサポートされますが、これにはtimefeature が必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
days: i32,
#[serde(with = "clickhouse::serde::time::date32")]
date: Date,
}DateTimeはu32またはそれをラップした newtype と相互変換でき、UNIX epoch からの経過秒数を表します。また、time::OffsetDateTimeもserde::time::datetimeを使うことでサポートされますが、これにはtimefeature が必要です。
#[derive(Row, Serialize, Deserialize)]
struct MyRow {
ts: u32,
#[serde(with = "clickhouse::serde::time::datetime")]
dt: OffsetDateTime,
}DateTime64(_)はi32またはそれをラップする newtype に相互変換され、Unix epoch からの経過時間を表します。また、time::OffsetDateTimeはserde::time::datetime64::*を使用することでサポートされますが、これにはtimefeature が必要です。
#[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,
}Variant、Dynamic、 (新しい)JSONデータ型はまだサポートされていません。
モック化
このクレートには、CHサーバーをモック化し、DDL、SELECT、INSERT、WATCH クエリをテストするためのユーティリティが用意されています。この機能は 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
}既知の制限事項
Variant、Dynamic、 (新しい)JSONデータ型は、まだサポートされていません。- サーバー側のパラメータバインドは、まだサポートされていません。追跡状況については、この issue を参照してください。
お問い合わせ
ご不明な点やサポートが必要な場合は、Community Slack または GitHub issues でお気軽にご連絡ください。