Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 35 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -219,8 +219,34 @@ try (Query q = db.borrowQuery()) {

### Cancel or Time Out a Query

`submit()` returns a `Completion`. `await(timeout, unit)` returns `false` if the query is still in flight; `cancel()`
stops it.
Give a query a timeout with `timeout(...)`, or every query a default with `query_timeout_ms` in the configuration
string. The timeout bounds the whole query, measured from `submit()`. When it expires the query is stopped, `await()`
throws a `QueryException` whose `isTimeout()` is `true`, and the pooled connection stays open for the next query.

```java
import io.questdb.client.QueryException;

try (Query q = db.borrowQuery()) {
q.sql("SELECT * FROM big_table ORDER BY ts")
.handler(handler)
.timeout(5, TimeUnit.SECONDS);
try {
q.submit().await();
} catch (QueryException e) {
if (!e.isTimeout()) {
throw e;
}
// the query ran longer than 5 seconds and was stopped
}
}
```

Servers that support per-query timeouts stop the query themselves; against older servers the client cancels it. If
the server does not end the query within `query_close_timeout_ms` of the timeout, `await()` throws anyway while the
connection finishes the aborted query in the background.

`submit()` returns a `Completion`. `await(timeout, unit)` only bounds the wait: it returns `false` while the query keeps
running. `cancel()` stops the query.

```java
import io.questdb.client.Completion;
Expand Down Expand Up @@ -540,6 +566,7 @@ schema::key1=value1;key2=value2;
| `query_pool_min` | `1` | Minimum query connections kept warm (`0` under `lazy_connect`) |
| `query_pool_max` | `4` | Maximum query connections |
| `acquire_timeout_ms` | `5000` | How long `borrowSender()`/`borrowQuery()` waits for a free slot |
| `query_close_timeout_ms` | `5000` | How long `Query.close()` waits for a running query; grace of `query_timeout_ms` |
| `idle_timeout_ms` | `60000` | How long a pooled connection may stay idle before it is reaped |
| `max_lifetime_ms` | `1800000` | Maximum lifetime of a pooled connection before it is recycled |

Expand All @@ -554,6 +581,12 @@ Applied by the query pool to select and fail over between the nodes in the `addr
| `zone` | | Prefer same-zone endpoints for `target=any`/`replica` (opaque, case-insensitive) |
| `client_id` | | Opaque client identifier surfaced server-side for observability |

### Query keys

| Key | Default | Description |
| ------------------ | ------- | ------------------------------------------------------------------------------------------ |
| `query_timeout_ms` | `0` | Default per-query timeout in milliseconds (`0` = none); override per query with `timeout()` |

The ingest side also accepts store-and-forward and reconnection tuning keys (`auto_flush_*`, `initial_connect_retry`,
`reconnect_*`, `request_durable_ack`, `sf_*`, `max_frame_rejections`, `poison_min_escalation_window_millis`, …).
`sf_max_total_bytes` caps the unacknowledged data all pooled senders buffer together — 128 MiB of memory by default,
Expand Down
4 changes: 4 additions & 0 deletions core/src/main/java/io/questdb/client/Completion.java
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,10 @@ public interface Completion {
/**
* Blocks up to the given timeout. Returns {@code true} if the query
* completed, {@code false} on timeout.
* <p>
* This bounds only how long the caller waits: on {@code false} the query
* keeps running. To stop queries that run too long, set a query timeout
* ({@link Query#timeout(long, TimeUnit)} or {@code query_timeout_ms}).
*
* @throws QueryException if the server reported an error or
* {@link #cancel()} won the race
Expand Down
44 changes: 41 additions & 3 deletions core/src/main/java/io/questdb/client/Query.java
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,10 @@
import io.questdb.client.cutlass.qwp.client.QwpBindSetter;
import io.questdb.client.cutlass.qwp.client.QwpColumnBatchHandler;
import io.questdb.client.cutlass.qwp.client.QwpServerInfo;
import io.questdb.client.cutlass.qwp.protocol.QwpConstants;

import java.io.Closeable;
import java.util.concurrent.TimeUnit;

/**
* A query handle leased from the {@link QuestDB} pool via
Expand All @@ -43,9 +45,9 @@
* creates one small lease handle per borrow (often scalar-replaced by the JIT
* when used with try-with-resources).
* <p>
* Lifecycle: configure with {@link #sql}, optional {@link #binds}, and
* {@link #handler}, then call {@link #submit()} to obtain a {@link Completion}
* and {@code await()} it before the next {@link #submit()}.
* Lifecycle: configure with {@link #sql}, optional {@link #binds} and
* {@link #timeout}, and {@link #handler}, then call {@link #submit()} to obtain
* a {@link Completion} and {@code await()} it before the next {@link #submit()}.
* <p>
* Thread safety: not thread-safe and single-flight -- one in-flight query per
* handle. To run queries concurrently, borrow one handle per concurrent query.
Expand Down Expand Up @@ -132,4 +134,40 @@ public interface Query extends Closeable {
* before the first successful bind
*/
QwpServerInfo serverInfo();

/**
* Sets the query timeout for subsequent {@link #submit()} calls on this
* handle, overriding the {@code query_timeout_ms} default from the
* connection string; {@code 0} runs queries without a timeout. It stays in
* effect until changed, and every borrow starts from the configured
* default.
* <p>
* The timeout is measured from {@code submit()} and bounds the whole query:
* server execution, any failover, and the time the handler spends in its
* callbacks (it is checked between result batches, so a handler blocked
* inside {@code onBatch} is not interrupted). When it expires, the query is
* stopped -- by the server when it supports per-query timeouts, otherwise by
* a client-side cancel -- no further result batch reaches the handler, the
* handler's {@code onError} receives
* {@link QwpConstants#STATUS_QUERY_TIMEOUT}, and {@link Completion#await()}
* throws a {@link QueryException} whose {@link QueryException#isTimeout()}
* is {@code true}. The pooled connection stays open and authenticated, and
* serves the next query.
* <p>
* Should the server not end the query within {@code query_close_timeout_ms}
* (see {@link QuestDBBuilder#queryCloseTimeoutMillis(long)}) of the timeout,
* the caller is released with the timeout anyway while the connection
* finishes draining the aborted query; only a connection that stays silent
* for that long once more is closed and replaced.
* <p>
* Unlike {@link Completion#await(long, TimeUnit)}, which only bounds how
* long the caller waits, the query timeout stops the query.
*
* @param timeout the timeout, {@code 0} for none; a positive value below
* one millisecond counts as one millisecond
* @param unit the unit of {@code timeout}
* @return this handle
* @throws IllegalArgumentException when {@code timeout} is negative
*/
Query timeout(long timeout, TimeUnit unit);
}
15 changes: 14 additions & 1 deletion core/src/main/java/io/questdb/client/QueryException.java
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@

package io.questdb.client;

import io.questdb.client.cutlass.qwp.protocol.QwpConstants;

/**
* Thrown from {@link Completion#await()} / {@link Completion#await(long, java.util.concurrent.TimeUnit)}
* when the server reported an error for the corresponding {@link Query},
Expand All @@ -32,7 +34,9 @@
* <p>
* The original wire-level status byte is exposed via {@link #getStatus()} so
* callers can distinguish cancellation from schema errors etc. without
* string-matching the message.
* string-matching the message. A query that ran past its timeout (see
* {@link Query#timeout(long, java.util.concurrent.TimeUnit)}) reports
* {@link QwpConstants#STATUS_QUERY_TIMEOUT}; {@link #isTimeout()} tests for it.
*/
public class QueryException extends RuntimeException {

Expand All @@ -56,4 +60,13 @@ public QueryException(byte status, String message, Throwable cause) {
public byte getStatus() {
return status;
}

/**
* Returns {@code true} when the query ran past its timeout, whether the
* server or the client detected it. The connection that ran it stays
* usable.
*/
public boolean isTimeout() {
return status == QwpConstants.STATUS_QUERY_TIMEOUT;
}
}
6 changes: 6 additions & 0 deletions core/src/main/java/io/questdb/client/QuestDBBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,12 @@ public QuestDBBuilder acquireTimeoutMillis(long millis) {
* letting the pool grow a fresh one. Bounds the close of a handle whose
* {@code submit()} is still running -- e.g. when the caller's own
* {@code await(timeout)} expired and they gave up. Defaults to 5000ms.
* <p>
* Also the grace period of the query timeout ({@code query_timeout_ms},
* {@link Query#timeout(long, java.util.concurrent.TimeUnit)}): how long a
* timed-out query may take to end before the caller is released with the
* timeout anyway, and how long its connection may then take to drain the
* aborted query before it is closed instead of reused.
*/
public QuestDBBuilder queryCloseTimeoutMillis(long millis) {
if (millis < 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,21 @@
* One event per {@code RESULT_BATCH} / {@code RESULT_END} / {@code QUERY_ERROR}
* received from the server, plus a synthetic error event if the connection drops
* mid-query.
* <p>
* The I/O thread stamps the events of a query with that query's id
* ({@link #requestId}), so {@code execute()} can tell its own events from the
* leftovers of an earlier query that its {@code execute()} stopped waiting for.
* Connection-level events -- the connection failed, the I/O thread stopped, the
* client closed -- carry {@link #ANY_REQUEST} and reach whichever query waits.
*/
public class QueryEvent {

/**
* {@link #requestId} of an event that belongs to no query in particular, such
* as a connection failure. Whichever query waits for events must see it.
*/
public static final long ANY_REQUEST = -1L;

public static final int KIND_BATCH = 0;
public static final int KIND_END = 1;
public static final int KIND_ERROR = 2;
Expand All @@ -50,19 +62,22 @@ public class QueryEvent {
public byte errorStatus; // valid for KIND_ERROR
public int kind;
public short opType; // valid for KIND_EXEC_DONE (matches CompiledQuery.SELECT/INSERT/etc.)
public long requestId = ANY_REQUEST; // query the event belongs to; see forRequest()
public long rowsAffected; // valid for KIND_EXEC_DONE
public long totalRows; // valid for KIND_END

public QueryEvent asBatch(QwpBatchBuffer buffer) {
this.kind = KIND_BATCH;
this.buffer = buffer;
this.requestId = ANY_REQUEST;
return this;
}

public QueryEvent asEnd(long totalRows) {
this.kind = KIND_END;
this.buffer = null;
this.totalRows = totalRows;
this.requestId = ANY_REQUEST;
return this;
}

Expand All @@ -71,6 +86,7 @@ public QueryEvent asError(byte status, String message) {
this.buffer = null;
this.errorStatus = status;
this.errorMessage = message;
this.requestId = ANY_REQUEST;
return this;
}

Expand All @@ -79,6 +95,7 @@ public QueryEvent asExecDone(short opType, long rowsAffected) {
this.buffer = null;
this.opType = opType;
this.rowsAffected = rowsAffected;
this.requestId = ANY_REQUEST;
return this;
}

Expand All @@ -87,9 +104,29 @@ public QueryEvent asTransportError(byte status, String message) {
this.buffer = null;
this.errorStatus = status;
this.errorMessage = message;
this.requestId = ANY_REQUEST;
return this;
}

/**
* Ties the event to the query with id {@code requestId}. Call it after the
* {@code asX()} builder, which leaves the event at {@link #ANY_REQUEST}.
*/
public QueryEvent forRequest(long requestId) {
this.requestId = requestId;
return this;
}

/**
* Returns {@code true} when the event belongs to a query other than
* {@code requestId}: a leftover of an earlier query that its
* {@code execute()} stopped waiting for. An {@link #ANY_REQUEST} event
* belongs to every query.
*/
public boolean isForOtherRequest(long requestId) {
return this.requestId != ANY_REQUEST && this.requestId != requestId;
}

/**
* Clears object references and resets primitive fields so a pooled event is
* safe to reuse across queries. The I/O thread calls the {@code asX(...)}
Expand All @@ -106,5 +143,6 @@ public void reset() {
this.opType = 0;
this.rowsAffected = 0;
this.totalRows = 0;
this.requestId = ANY_REQUEST;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,12 @@
* <strong>Exception contract:</strong> if any callback method throws, the
* exception propagates out of the {@link QwpQueryClient#execute} call on the
* caller's thread and no further callbacks fire for that query. The connection
* remains usable for subsequent queries.
* remains usable for subsequent queries: when {@link #onBatch} throws while the
* query is still running, {@code execute} cancels the query, and the next query
* skips whatever the cancelled one still sends. When {@link #onError} reports a
* query timeout at the end of the grace period, the aborted query is still
* running as well; if that callback throws, {@code execute} first drains the
* query, as it does when the callback returns, and only then rethrows.
*/
public interface QwpColumnBatchHandler {

Expand Down
Loading
Loading