questdb-rs 7.0.0

QuestDB Client Library for Rust
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
# Fast Ingestion of Data into QuestDB

The `ingress` module implements QuestDB's variant of the
[InfluxDB Line Protocol](https://questdb.io/docs/reference/api/ilp/overview/)
(ILP).

To get started:

* Use [`Sender::from_conf()`] to get the [`Sender`] object
* Populate a [`Buffer`] with one or more rows of data
* Send the buffer using [`sender.flush()`]Sender::flush

```rust no_run
use questdb::{
    Result,
    ingress::{
        Sender,
        Buffer,
        TimestampNanos}};
fn main() -> Result<()> {
   let mut sender = Sender::from_conf("http::addr=localhost:9000;")?;
  let mut buffer = sender.new_buffer();
   buffer
       .table("trades")?
       .symbol("symbol", "ETH-USD")?
       .symbol("side", "sell")?
       .column_f64("price", 2615.54)?
       .column_f64("amount", 0.00044)?
       .at(TimestampNanos::now())?;
   sender.flush(&mut buffer)?;
   Ok(())
}
```

# Configuration String

The easiest way to configure all the available parameters on a line sender is
the configuration string. The general structure is:

```plain
<transport>::addr=host:port;param1=val1;param2=val2;...
```

`transport` can be `http`, `https`, `tcp`, or `tcps`. See the full details on
supported parameters in a dedicated section below.

# Don't Forget to Flush

The sender and buffer objects are entirely decoupled. This means that the sender
won't get access to the data in the buffer until you explicitly call
[`sender.flush(&mut buffer)`](Sender::flush) or a variant. This may lead to a
pitfall where you drop a buffer that still has some data in it, resulting in
permanent data loss.

A common technique is to flush periodically on a timer and/or once the buffer
exceeds a certain size. You can check the buffer's size by the calling
[`buffer.len()`](Buffer::len).

The default `flush()` method clears the buffer after sending its data. If you
want to preserve its contents (for example, to send the same data to multiple
QuestDB instances), call
[`sender.flush_and_keep(&mut buffer)`](Sender::flush_and_keep) instead.

# Error Handling

The supported transport modes handle errors very differently. In a nutshell,
HTTP is much better at error handling.

## TCP

TCP doesn't report errors at all to the sender; instead, the server quietly
disconnects and you'll have to inspect the server logs to get more information
on the reason. When this has happened, the sender transitions into an error
state, and it is permanently unusable. You must drop it and create a new sender.
You can inspect the sender's error state by calling
[`sender.must_close()`](Sender::must_close).

## QWP/UDP

QWP/UDP is a best-effort datagram transport. A `flush()` call sends one or more
UDP datagrams and can report only local socket errors. A successful return does
not guarantee delivery, server-side processing, or flush-level atomicity.

When one logical flush spans multiple datagrams, some datagrams may already
have been emitted before a later send fails. In that case, retrying the same
buffer may duplicate rows that were already sent.

## QWP/WebSocket

QWP/WebSocket is a reliable, in-order transport. Each published frame is
assigned a frame sequence number (FSN) and tracked through a publication
lifecycle until the server acknowledges or rejects it. Transient socket
errors and reconnects are absorbed by the driver and do not surface to the
caller; only server rejections and terminal protocol violations are
reported.

Server rejections are classified as `Retriable`, `RetriableOther`, or
`Terminal`. Retriable rejections (e.g. write error, internal error, unknown
future status) reconnect and replay from the local store-and-forward
watermark. `RetriableOther` is reserved for server states such as
`NOT_WRITABLE`; it uses the same replay path but immediately rotates away
from the read-only endpoint when another configured endpoint is available.
Structured diagnostics are available with
[`sender.poll_qwp_ws_error()`](Sender::poll_qwp_ws_error), or through a
callback installed with
[`SenderBuilder::qwp_ws_error_handler`](SenderBuilder::qwp_ws_error_handler).

Terminal rejections (e.g. schema mismatch, parse error, security error) and
terminal WebSocket protocol violations latch the sender into a
permanently-unusable state. The next `flush()` call returns
  `ErrorCode::ServerRejection` carrying the structured diagnostic, and
  [`sender.must_close()`](Sender::must_close) returns `true`. The sender
  must be dropped and a new one constructed. The buffer passed to that
  final `flush()` is left unmodified, so its contents can be recovered
  before dropping the sender.

`flush()` returns once the frame has been appended to the local publication
log; the FSN is also available via
[`flush_and_get_fsn()`](Sender::flush_and_get_fsn). To wait for the server
to acknowledge every frame published so far, call
[`sender.wait()`](Sender::wait) with the desired
[`AckLevel`](crate::ingress::AckLevel) — the row-major counterpart to the
column-major store-and-forward `wait`.
[`sender.published_fsn()`](Sender::published_fsn) and
[`sender.acked_fsn()`](Sender::acked_fsn) provide non-blocking polls.

Configure `sf_dir` to recover the local publication log after reconnects and
producer-process restarts. The default `sf_durability=memory` mode relies on
the OS page cache and does not cover host power loss.
`sf_durability=periodic` checkpoints published mmap data in the background;
`sf_sync_interval_millis` defaults to `5000` in that mode. The interval is a
target cadence rather than a hard recovery-window guarantee because runner
scheduling and storage-sync latency add to it. At a segment boundary,
publication may receive normal store-and-forward backpressure until a
checkpoint makes the segment safe to rotate.

Periodic durability protects the local replay log. It is independent of the
QuestDB Enterprise server-side durable ACK barrier selected by
`request_durable_ack=on`; configure both when end-to-end durability is
required.

In `manual` progress mode no background thread observes the transport.
Server-side state — including terminal diagnostics — only becomes visible when the user
calls [`sender.drive_once()`](Sender::drive_once) or any sender method
that drives the transport (such as `flush` or `wait`). As a
consequence, `must_close()` on an otherwise-idle manual sender does not
reflect a terminal diagnostic until the next drive.

## HTTP

HTTP distinguishes between recoverable and non-recoverable errors. For
recoverable ones, it enters a retry loop with exponential backoff, and reports
the error to the caller only after it has exhausted the retry time budget
(configuration parameter: `retry_timeout`).

`sender.flush()` and variant methods communicate the error in the `Result`
return value. The category of the error is signalled through the
[`ErrorCode`](crate::error::ErrorCode) enum, and it's accompanied with an error
message.

After the sender has signalled an error, it remains usable. You can handle the
error as appropriate and continue using it.

# Health Check

The QuestDB server has a "ping" endpoint you can access to see if it's alive,
and confirm the version of the InfluxDB that it is compatible with at a protocol
level.

```shell
curl -I http://localhost:9000/ping
```

Example of the expected response:

```plain
HTTP/1.1 204 OK
Server: questDB/1.0
Date: Fri, 2 Feb 2024 17:09:38 GMT
Transfer-Encoding: chunked
Content-Type: text/plain; charset=utf-8
X-Influxdb-Version: v2.7.4
```

# Configuration Parameters

In the examples below, we'll use configuration strings. We also provide the
[`SenderBuilder`] to programmatically configure the sender. The methods on
[`SenderBuilder`] match one-for-one with the keys in the configuration string.

## Authentication

To establish an
[authenticated](https://questdb.io/docs/reference/api/ilp/overview/#authentication)
and TLS-encrypted connection, use the `https` or `tcps` protocol, and use the
configuration options appropriate for the authentication method.

Here are quick examples of configuration strings for each authentication method
we support:

### HTTP Token Bearer Authentication

```no_run
# use questdb::{Result, ingress::Sender};
# fn main() -> Result<()> {
let mut sender = Sender::from_conf(
    "https::addr=localhost:9000;token=Yfym3fgMv0B9;"
)?;
# Ok(())
# }
```

* `token`: the authentication token

### HTTP Basic Authentication

```no_run
# use questdb::{Result, ingress::Sender};
# fn main() -> Result<()> {
let mut sender = Sender::from_conf(
    "https::addr=localhost:9000;username=testUser1;password=Yfym3fgMv0B9;"
)?;
# Ok(())
# }
```

* `username`: the username
* `password`: the password

### TCP Elliptic Curve Digital Signature Algorithm (ECDSA)

```no_run
# use questdb::{Result, ingress::Sender};
# fn main() -> Result<()> {
let mut sender = Sender::from_conf(
    "tcps::addr=localhost:9009;username=testUser1;token=5UjEA0;token_x=fLKYa9;token_y=bS1dEfy;"
)?;
# Ok(())
# }
```

The four ECDSA components are:

* `username`, aka. _kid_
* `token`, aka. _d_
* `token_x`, aka. _x_
* `token_y`, aka. _y_

### Authentication Timeout

You can specify how long the client should wait for the authentication request
to resolve. The configuration parameter is:

* `auth_timeout` (milliseconds, default 15 seconds)

For QWP/WebSocket configuration strings, the Java-compatible spelling is also
accepted:

* `auth_timeout_ms` (milliseconds, default 15 seconds)

## Encryption on the Wire: TLS

To enable TLS on the QuestDB Enterprise server, refer to the [QuestDB Enterprise
TLS documentation](https://questdb.io/docs/operations/tls/).

*Note*: QuestDB Open Source does not support TLS natively. To use TLS with
QuestDB Open Source, use a TLS proxy such as
[HAProxy](http://www.haproxy.org/).

We support several certification authorities (sources of PKI root certificates).
To select one, use the `tls_ca` config option. These are the supported variants:

* `tls_ca=webpki_roots;` use the roots provided in the standard Rust crate
  [webpki-roots]https://crates.io/crates/webpki-roots

* `tls_ca=os_roots;` use the OS-provided certificate store

* `tls_ca=webpki_and_os_roots;` combine both of the above

* `tls_roots=/path/to/root-ca.pem;` get the root certificates from the specified
  file. Main purpose is for testing with self-signed certificates. _Note:_ this
  automatically sets `tls_ca=pem_file`.

* `tls_roots_password=<secret>;` unlocks a JKS / PKCS#12 keystore named by
  `tls_roots`. QWP/WebSocket (`wss::`) **only** — ILP/TCP and ILP/HTTP read
  unencrypted PEM via rustls and reject this key. With the password set, the
  `tls_roots` file is interpreted as a Java KeyStore (auto-detected: JKS magic
  `0xFEEDFEED`, or PKCS#12 ASN.1 SEQUENCE) and its trusted-certificate entries
  become the rustls root store. Mirrors the Java reference client's
  `tls_roots_password` connect-string key.

See our notes on [how to generate a self-signed
certificate](https://github.com/questdb/c-questdb-client/tree/main/tls_certs).

* `tls_verify=unsafe_off;` tells the QuestDB client to ignore all CA roots and
  accept any server certificate without checking. You can use it as a last
  resort, when you weren't able to apply the above approach with a self-signed
  certificate. You should **never use it in production** as it defeats security
  and allows a man-in-the middle attack.

## HTTP Timeouts

Instead of a fixed timeout value, we use a flexible timeout that depends on the
size of the HTTP request payload (how much data is in the buffer that you're
flushing). You can configure it using two options:

* `request_min_throughput` (bytes per second, default 100 KiB/s): divide the
  payload size by this number to determine for how long to keep sending the
  payload before timing out.
* `request_timeout` (milliseconds, default 10 seconds): additional time
  allowance to account for the fixed latency of the request-response roundtrip.

Finally, the client will keep retrying the request if it experiences errors. You
can configure the total time budget for retrying:

* `retry_timeout` (milliseconds, default 10 seconds)

# Usage Considerations

## Transactional Flush

When using HTTP, you can arrange that each `flush()` call happens within its own
transaction. For this to work, your buffer must contain data that targets only
one table. This is because QuestDB doesn't support multi-table transactions.

In order to ensure in advance that a flush will be transactional, call
[`sender.flush_and_keep_with_flags(&mut buffer, true)`](Sender::flush_and_keep_with_flags).
This call will refuse to flush a buffer if the flush wouldn't be transactional.

## When to Choose the TCP Transport?

As discussed above, the TCP transport mode is raw and simplistic: it doesn't
report any errors to the caller (the server just disconnects), has no automatic
retries, requires manual handling of connection failures, and doesn't support
transactional flushing.

However, TCP has a lower overhead than HTTP and it's worthwhile to try out as an
alternative in a scenario where you have a constantly high data rate and/or deal
with a high-latency network connection.

## Array Datatype

The [`Buffer::column_arr`](Buffer::column_arr) method supports efficient ingestion of N-dimensional
arrays using several convenient types:

- native Rust arrays and slices (up to 3-dimensional)
- native Rust vectors (up to 3-dimensional)
- arrays from the [`ndarray`]https://docs.rs/ndarray crate

You must use protocol version 2 to ingest arrays. The HTTP transport will
automatically enable it as long as you're connecting to an up-to-date QuestDB
server (version 9.0.0 or later), but with TCP you must explicitly specify it in
the configuration string: `protocol_version=2;`.

**Note**: QuestDB server version 9.0.0 or later is required for array support.

## Timestamp Column Name

The InfluxDB Line Protocol (ILP) does not give a name to the designated timestamp,
so if you let this client auto-create the table, it will have the default `timestamp` name.
To use a custom name, say `my_ts`, pre-create the table with the desired
timestamp column name:

To address this, issue a `CREATE TABLE` statement to create the table in advance.
Note the `timestamp(my_ts)` clause at the end specifies the designated timestamp.

```sql
CREATE TABLE IF NOT EXISTS 'trades' (
  symbol SYMBOL capacity 256 CACHE,
  side SYMBOL capacity 256 CACHE,
  price DOUBLE,
  amount DOUBLE,
  my_ts TIMESTAMP
) timestamp (my_ts) PARTITION BY DAY WAL;
```

You can use the `CREATE TABLE IF NOT EXISTS` construct to make sure the table is
created, but without raising an error if the table already exists.

## Sequential Coupling in the Buffer API

The fluent API of [`Buffer`] has sequential coupling: there's a certain order in
which you are expected to call the methods. For example, you must write the
symbols before the columns, and you must terminate each row by calling either
[`at`](Buffer::at) or [`at_now`](Buffer::at_now). Refer to the [`Buffer`] doc
for the full rules and a flowchart.

## Optimization: Avoid Revalidating Names

The client validates every name you provide. To avoid the redundant CPU work of
re-validating the same names on every row, create pre-validated [`ColumnName`]
and [`TableName`] values:

```no_run
# use questdb::Result;
use questdb::ingress::{
    TableName,
    ColumnName,
    Buffer,
    SenderBuilder,
    TimestampNanos};
# fn main() -> Result<()> {
let mut sender = SenderBuilder::from_conf("https::addr=localhost:9000;")?.build()?;
let mut buffer = sender.new_buffer();
let table_name = TableName::new("trades")?;
let price_name = ColumnName::new("price")?;
buffer.table(table_name)?.column_f64(price_name, 2615.54)?.at(TimestampNanos::now())?;
buffer.table(table_name)?.column_f64(price_name, 39269.98)?.at(TimestampNanos::now())?;
# Ok(())
# }
```

## Handling Optional Data (NULLs)

In QuestDB, `NULL` values are represented by simply omitting the column for that specific row.

To make working with Rust's `Option<T>` ergonomic and keep the fluent builder chain unbroken, the [`Buffer`] API provides `_opt` variants for all column methods (e.g., [`column_str_opt`](Buffer::column_str_opt), [`column_f64_opt`](Buffer::column_f64_opt)).

If the provided value is `Some(v)`, the column is written normally. If the value is `None`, the method acts as a no-op and skips the column.

**Note on ownership:** For types that implement `Copy` (like `i64`, `f64`, `bool`), you can pass the `Option` directly. For heap-allocated types like `String` or `Vec`, use `.as_ref()` or `.as_deref()` to pass a reference without consuming the original value.

```rust no_run
# use questdb::Result;
use questdb::ingress::{Buffer, SenderBuilder, TimestampNanos};

# fn main() -> Result<()> {
  let mut sender = SenderBuilder::from_conf("https::addr=localhost:9000;")?.build()?;
  let mut buffer = sender.new_buffer();
  struct SensorData {
      sensor_id: Option<String>,
      temperature: Option<f64>,
      humidity: Option<f64>,
  }

  let data = SensorData { sensor_id: Some("sensor-1".to_string()), temperature: Some(22.5), humidity: None };

  buffer
      .table("sensors")?
      .symbol("location", "factory-1")?
      // Writes the sensor_id column if it's Some, otherwise skips it (stored as NULL in QuestDB)
      .symbol_opt("sensor_id", data.sensor_id.as_deref())?
      // Writes the temperature column
      .column_f64_opt("temperature", data.temperature)?
      // Silently skips the humidity column (stored as NULL in QuestDB)
      .column_f64_opt("humidity", data.humidity)?
      .at(TimestampNanos::now())?;
#   Ok(())
# }
```

## Decimal Datatype

The [`Buffer::column_dec`](Buffer::column_dec) method supports efficient ingestion of decimal values using several convenient types:

- native Rust String slices
- decimals from the [`rust_decimal`]https://docs.rs/rust_decimal crate
- decimals from the [`bigdecimal`]https://docs.rs/bigdecimal crate

You must use protocol version 3 to ingest decimals. The HTTP transport will
automatically enable it as long as you're connecting to an up-to-date QuestDB
server (version 9.2.0 or later), but with TCP you must explicitly specify it in
the configuration string: `protocol_version=3;`.

**Note**: QuestDB server version 9.2.0 or later is required for decimal support.

## Check out the CONSIDERATIONS Document

The [Library
considerations](https://github.com/questdb/c-questdb-client/blob/main/doc/CONSIDERATIONS.md)
document covers these topics:

* Threading
* Differences between the InfluxDB Line Protocol and QuestDB Data Types
* Data Quality
* Client-side checks and server errors
* Flushing
* Disconnections, data errors and troubleshooting

# Troubleshooting Common Issues

## Infrequent Flushing

If the data doesn't appear in the database in a timely manner, you may not be
calling [`flush()`](Sender::flush) often enough.

## Debug disconnects and inspect errors

If you're using ILP-over-TCP, it doesn't report any errors to the client.
Instead, on error, the server terminates the connection, and logs any error
messages in [server logs](https://questdb.io/docs/troubleshooting/log/).

To inspect or log a buffer's contents before you send it, call
[`buffer.as_bytes()`](Buffer::as_bytes).

This byte-level inspection is only meaningful for ILP buffers. QWP buffers are
encoded into UDP datagrams during [`flush()`](Sender::flush), so
[`buffer.as_bytes()`](Buffer::as_bytes) is not useful there.