ClickHouse 연결을 위한 공식 Rust 클라이언트입니다. 최초 개발자는 Paul Loyd이며, 클라이언트 소스 코드는 GitHub 리포지토리에서 확인할 수 있습니다.
개요
- 행 인코딩/디코딩에
serde를 사용합니다. serde속성인skip_serializing,skip_deserializing,rename을 지원합니다.- HTTP 전송을 통해
RowBinary포맷을 사용합니다.- 향후 TCP를 통해
Native를 사용하도록 전환할 계획입니다.
- 향후 TCP를 통해
- TLS(
native-tls및rustls-tls기능 사용)를 지원합니다. - 압축 및 압축 해제(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(_)variant를 활성화합니다. 활성화되면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 버전 호환성
이 클라이언트는 LTS 및 그 이후 버전의 ClickHouse는 물론 ClickHouse Cloud와도 호환됩니다.
v22.6 이전의 ClickHouse 서버는 드물게 일부 경우에 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"));관련 항목:
- 클라이언트 리포지토리의 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미만인 경우에만 배치를 원자적으로 삽입합니다.
Async insert (서버 측 배칭)
수신 데이터에 대해 클라이언트 측 배칭을 하지 않으려면 ClickHouse asynchronous inserts를 사용할 수 있습니다. 이렇게 하려면 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");관련 항목:
- 클라이언트 리포지토리의 async insert 예시
Inserter 기능(클라이언트 측 배칭)
inserter 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?;Inserter는 임계값(max_bytes,max_rows,period) 중 하나에 도달하면commit()에서 진행 중인 삽입을 종료합니다.- 진행 중인
INSERT종료 간격에는with_period_bias로 바이어스를 적용할 수 있으며, 이를 통해 병렬 inserter로 인한 부하 급증을 방지할 수 있습니다. Inserter::time_left()는 현재 주기가 언제 끝나는지 감지하는 데 사용할 수 있습니다. 스트림에서 항목이 드물게 들어오는 경우, 한도를 확인하려면Inserter::commit()을 다시 호출하십시오.- 시간 임계값은
inserter의 성능을 높이기 위해 quanta 크레이트를 사용해 구현됩니다.test-util이 활성화된 경우에는 사용되지 않으므로, 사용자 정의 테스트에서는tokio::time::advance()로 시간을 관리할 수 있습니다. 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 query log에서 쿼리를 식별할 수 있습니다.
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 예시.
세션 ID
query_id와 마찬가지로 session_id를 설정하면 동일한 session에서 SQL 문을 실행할 수 있습니다. 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 헤더 예시를 참조하십시오.
사용자 지정 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)에 매핑되거나, 그 반대로도 매핑될 수 있습니다. 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로 또는 그 반대로 매핑됩니다.uuid기능이 필요합니다.
#[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기능이 필요하며,serde::time::date32를 사용하면time::Date도 지원합니다.
#[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을 사용해 지원되며, 이 경우time기능이 필요합니다.
#[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::*를 사용해 지원되며, 이 경우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::*helpers를 사용할 때는::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 기능을 활성화하면 사용할 수 있습니다. 반드시 개발 의존성으로만 사용하세요.
예시를 참조하세요.
문제 해결
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데이터 타입은 아직 지원되지 않습니다.- 서버 측 매개변수 바인딩은 아직 지원되지 않습니다. 관련 진행 상황은 this issue에서 추적하고 있습니다.
문의하기
질문이 있거나 도움이 필요하시면 Community Slack 또는 GitHub 이슈를 통해 언제든지 문의하세요.