عميل Rust الرسمي للاتصال بـ ClickHouse، وقد طوّره أساسًا Paul Loyd. الشيفرة المصدرية للعميل متاحة في مستودع GitHub.
نظرة عامة
- يستخدم
serdeلترميز الصفوف وفك ترميزها. - يدعم سمات
serde:skip_serializingوskip_deserializingوrename. - يستخدم تنسيق
RowBinaryعبر نقل HTTP.- توجد خطط للانتقال إلى
Nativeعبر TCP.
- توجد خطط للانتقال إلى
- يدعم TLS (عبر ميزتَي
native-tlsوrustls-tls). - يدعم الضغط وفك الضغط (LZ4).
- يوفّر واجهات برمجة تطبيقات للاستعلام عن البيانات أو إدراجها، وتنفيذ أوامر 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— تضيف mocks. راجع المثال. استخدمها فقط فيdev-dependencies.watch— تفعّل وظيفةclient.watch. راجع القسم المقابل لمزيد من التفاصيل.uuid— تضيفserde::uuidللعمل مع حزمة uuid.time— تضيفserde::timeللعمل مع حزمة time.
توافق إصدارات ClickHouse
العميل متوافق مع إصدارات LTS أو الأحدث من ClickHouse، وكذلك مع ClickHouse Cloud.
يتعامل خادم ClickHouse الأقدم من 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 في مستودع عميل. ينبغي أن ينطبق هذا أيضًا على اتصالات 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 الدُفعات بصورة ذرّية فقط إذا كانت جميع الصفوف تقع ضمن partition نفسها وكان عددها أقل من
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 (التجميع على جهة العميل)
يتطلب ذلك تفعيل ميزة 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عملية الإدراج النشطة فيcommit()إذا تم بلوغ أيٍّ من العتبات (max_bytes،max_rows،period). - يمكن إزاحة الفاصل الزمني بين إنهاء أوامر
INSERTالنشطة باستخدامwith_period_biasلتجنّب ارتفاعات الحمل الناتجة عن أدوات الإدراج المتوازية. - يمكن استخدام
Inserter::time_left()لاكتشاف موعد انتهاء الفترة الحالية. استدعِInserter::commit()مرة أخرى للتحقق من الحدود إذا كان التدفق يُصدر العناصر بوتيرة متباعدة. - تُنفَّذ العتبات الزمنية باستخدام حزمة quanta لتسريع
inserter. ولا يُستخدم ذلك إذا كانtest-utilمُمكّنًا (وبالتالي يمكن التحكم في الوقت عبرtokio::time::advance()في الاختبارات المخصّصة). - تُدرَج جميع الصفوف بين استدعاءات
commit()ضمن عبارةINSERTنفسها.
تنفيذ DDLs
في النشر أحادي العقدة، يكفي تنفيذ DDLs كما يلي:
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 لضبط إعدادات عامة لجميع الاستعلامات.
معرّف الاستعلام
باستخدام .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 مخصصة
إذا كنت تستخدم المصادقة عبر proxy أو تحتاج إلى تمرير رؤوس مخصصة، فيمكنك القيام بذلك كما يلي:
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)أوnewtypesالمبنية عليها. - لا يتوفر دعم مباشر لـ
(U)Int256، ولكن يوجد حل بديل لذلك. - يقابل
Float(32|64)في التحويل من/إلىf(32|64)المناظرة أوnewtypesالمبنية عليها. - يقابل
Decimal(32|64|128)في التحويل من/إلىi(32|64|128)المناظرة أوnewtypesالمبنية عليها. ويكون استخدامfixnumأو أي تنفيذ آخر للأعداد العشرية الثابتة ذات الإشارة أكثر ملاءمة. - يقابل
Booleanفي التحويل من/إلىboolأوnewtypesالمبنية عليه. - يقابل
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(_)يقابل ذهابًا وإيابًا أيslice، مثل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وJSON(الجديدة) غير مدعومة حتى الآن. - ربط المعلّمات على جهة الخادم غير مدعوم حتى الآن؛ راجع هذه التذكرة لمتابعة الحالة.
تواصل معنا
إذا كانت لديك أي أسئلة أو كنت بحاجة إلى مساعدة، فلا تتردد في التواصل معنا عبر Slack الخاص بالمجتمع أو من خلال تذاكر GitHub.