Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

عميل Rust لـ ClickHouse

عميل 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"));

انظر أيضًا:

تحديد الصفوف

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

راجع أيضًا:

ميزة 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.

Navigation