Skip to main content

ClickHouseConfig

Struct ClickHouseConfig 

Source
pub struct ClickHouseConfig {
Show 18 fields pub url: String, pub username: Option<String>, pub password: Option<String>, pub database: Option<String>, pub table: String, pub columns: Option<BTreeMap<String, String>>, pub async_insert: bool, pub wait_for_async_insert: Option<bool>, pub cursor_column: Option<String>, pub cursor_id: Option<String>, pub checkpoint_store: Option<String>, pub select_columns: Option<String>, pub polling_interval_ms: Option<u64>, pub max_polling_interval_ms: Option<u64>, pub request_timeout_ms: Option<u64>, pub connect_timeout_ms: Option<u64>, pub tls: TlsConfig, pub compression: Compression,
}
Expand description

ClickHouse endpoint configuration (talks the ClickHouse HTTP interface).

As a publisher it batch-inserts messages using FORMAT JSONEachRow — by default the whole message payload (which must be a JSON object) becomes one row; set columns to build each row from explicit ${payload:<field>} / ${metadata:<key>} tokens instead. As a consumer it reads an existing table non-destructively by paging over a monotonic cursor_column (ClickHouse has no native queue/pub-sub), serializing each row to a JSON payload.

Fields§

§url: String

ClickHouse HTTP endpoint URL, e.g. http://localhost:8123 (or https://…). If it contains userinfo, it will be treated as a secret.

§username: Option<String>

Optional username. Takes precedence over any credentials embedded in the url. Defaults to default.

§password: Option<String>

Optional password. Takes precedence over any credentials embedded in the url.

§database: Option<String>

Database name. Defaults to default.

§table: String

The table to read from / write to. May be schema-qualified (db.table).

§columns: Option<BTreeMap<String, String>>

(Publisher only) Optional per-column mapping. Each entry maps a target column name to a value token: ${payload:<field>} takes the top-level JSON field <field> of the payload (JSON type preserved), ${metadata:<key>} takes message.metadata["<key>"] (as a string), and any other value is inserted literally. When omitted, the whole payload JSON object is inserted as one row.

§async_insert: bool

(Publisher only) If true, set the ClickHouse async_insert=1 server setting so inserts are buffered server-side. Defaults to false.

§wait_for_async_insert: Option<bool>

(Publisher only) With async_insert, wait for the server to flush before acking. Defaults to true (durable). False = fire-and-forget: faster, but a crash before flush can drop the batch.

§cursor_column: Option<String>

(Consumer only) Read an existing table non-destructively and resumably, paging by this monotonic column (SELECT … WHERE {cursor_column} > {last} ORDER BY {cursor_column} ASC LIMIT n) and persisting the last read value under cursor_id.

§cursor_id: Option<String>

(Consumer only) Cursor id used to key the persisted resume position. Without it, progress is not persisted and every restart re-copies from the beginning.

§checkpoint_store: Option<String>

(Consumer only) Where to persist the resume cursor. Because ClickHouse is unsuited to per-row cursor upserts, a durable checkpoint requires an external store URL:

  • file:///var/lib/mqb/cursors.json → local JSON file
  • postgres://user@host/db/table / mysql://host/db/table → external SQL table (table optional)
  • mongodb://host/db/collection → external MongoDB collection (collection optional)
  • s3://bucket/prefix (also gs://, az://, abfs://) → cloud object store; creds via env

May embed connection credentials, so it is treated as a secret.

§select_columns: Option<String>

(Consumer only) Columns to select in cursor_column mode. Defaults to *.

§polling_interval_ms: Option<u64>

(Consumer only) Polling interval in milliseconds when the table is drained. Defaults to 100ms.

§max_polling_interval_ms: Option<u64>

(Consumer only) If set, the poll interval backs off exponentially from polling_interval_ms up to this value while drained, resetting on new rows. Unset = constant interval.

§request_timeout_ms: Option<u64>

Request timeout in milliseconds for ClickHouse HTTP calls (inserts, cursor reads, status). Unset = no timeout (wait indefinitely), which suits very large batch inserts.

§connect_timeout_ms: Option<u64>

Connection (TCP + TLS handshake) timeout in milliseconds. Defaults to 10000ms.

§tls: TlsConfig

TLS configuration for https:// connections.

§compression: Compression

HTTP body compression for inserts and cursor reads (none, gzip, lz4, zstd). Applied as Content-Encoding on the request body and negotiated on the response via Accept-Encoding. lz4/zstd are faster than gzip; all are understood natively by ClickHouse. Defaults to gzip.

Trait Implementations§

Source§

impl Clone for ClickHouseConfig

Source§

fn clone(&self) -> ClickHouseConfig

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ClickHouseConfig

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Default for ClickHouseConfig

Source§

fn default() -> ClickHouseConfig

Returns the “default value” for a type. Read more
Source§

impl<'de> Deserialize<'de> for ClickHouseConfig

Source§

fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>
where __D: Deserializer<'de>,

Deserialize this value from the given Serde deserializer. Read more
Source§

impl SecretExtractor for ClickHouseConfig

Source§

fn extract_secrets( &mut self, prefix: &str, secrets: &mut HashMap<String, String>, )

Extracts secrets into the provided map using the given prefix, and clears them from self.
Source§

impl Serialize for ClickHouseConfig

Source§

fn serialize<__S>(&self, __serializer: __S) -> Result<__S::Ok, __S::Error>
where __S: Serializer,

Serialize this value into the given Serde serializer. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DeserializeOwned for T
where T: for<'de> Deserialize<'de>,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more