Skip to content
Merged
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
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -556,6 +556,9 @@ Applied by the query pool to select and fail over between the nodes in the `addr

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,
10 GiB of disk with `sf_dir` — so a larger `sender_pool_max` adds connections, not buffer memory. The only overshoot
is the minimum working set every live sender keeps: two segments (`2 × sf_max_segment_bytes`, 8 MiB by default).
`sf_durability=periodic` checkpoints mmap-published data in the background; `sf_sync_interval_millis` defaults to `5000`
in that mode. The interval is a target cadence: JVM scheduling and storage-sync latency add to the actual loss window.
Use `request_durable_ack=on` when end-to-end server durability is also required. See the
Expand Down
5 changes: 5 additions & 0 deletions core/src/main/java/io/questdb/client/QuestDBBuilder.java
Original file line number Diff line number Diff line change
Expand Up @@ -400,6 +400,11 @@ public QuestDBBuilder queryPoolSize(int size) {

/**
* Maximum sender-pool size. Defaults to 4.
* <p>
* The pooled senders share one {@code sf_max_total_bytes} budget for data
* the server has not acknowledged yet, so a larger pool adds connections,
* not buffer memory -- apart from the minimum working set of two segments
* every live sender keeps.
*/
public QuestDBBuilder senderPoolMax(int max) {
if (max < 1) {
Expand Down
126 changes: 98 additions & 28 deletions core/src/main/java/io/questdb/client/Sender.java
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@
import io.questdb.client.cutlass.qwp.client.sf.cursor.CursorWebSocketSendLoop;
import io.questdb.client.cutlass.qwp.client.sf.cursor.OrphanScanner;
import io.questdb.client.cutlass.qwp.client.sf.cursor.PersistedSymbolDict;
import io.questdb.client.cutlass.qwp.client.sf.cursor.SegmentBudget;
import io.questdb.client.cutlass.qwp.client.sf.cursor.SlotLock;
import io.questdb.client.cutlass.qwp.client.sf.cursor.MmapSegmentCorruptionException;
import io.questdb.client.cutlass.qwp.client.sf.cursor.SfRecoveryException;
Expand Down Expand Up @@ -1162,6 +1163,9 @@ public int getConnectTimeout() {
private SfDurability sfDurability = SfDurability.MEMORY;
private long sfMaxSegmentBytes = PARAMETER_NOT_SET_EXPLICITLY;
private long sfMaxTotalBytes = PARAMETER_NOT_SET_EXPLICITLY;
// Budget shared with other senders (a sender pool's), or null for a
// private budget of sfMaxTotalBytes. See storeAndForwardSharedBudget.
private SegmentBudget sfSharedBudget;
private long sfSyncIntervalMillis = PARAMETER_NOT_SET_EXPLICITLY;
private boolean shouldDestroyPrivKey;
private boolean tlsEnabled;
Expand Down Expand Up @@ -1484,18 +1488,8 @@ public Sender build() {
// (same lock-free architecture, no disk involvement).
// Durability-combination validation lives in validateParameters
// so build() and no-connect validation apply the same rules.
long actualSfMaxSegmentBytes = sfMaxSegmentBytes == PARAMETER_NOT_SET_EXPLICITLY
? DEFAULT_SEGMENT_BYTES
: sfMaxSegmentBytes;
// Default cap depends on backing: RAM (memory mode) is tight
// by default; disk (SF mode) is cheap so the default is
// generous enough that normal traffic never hits it.
long defaultMaxTotal = sfDir == null
? DEFAULT_MAX_BYTES_MEMORY
: DEFAULT_MAX_BYTES_SF;
long actualSfMaxTotalBytes = sfMaxTotalBytes == PARAMETER_NOT_SET_EXPLICITLY
? Math.max(defaultMaxTotal, actualSfMaxSegmentBytes * 2)
: sfMaxTotalBytes;
long actualSfMaxSegmentBytes = resolveSfMaxSegmentBytes();
long actualSfMaxTotalBytes = resolveSfMaxTotalBytes();
long actualCloseFlushTimeoutMillis = closeFlushTimeoutMillis == CLOSE_FLUSH_TIMEOUT_NOT_SET
? DEFAULT_CLOSE_FLUSH_TIMEOUT_MILLIS
: closeFlushTimeoutMillis;
Expand Down Expand Up @@ -1624,9 +1618,9 @@ public Sender build() {
CursorSendEngine cursorEngine;
try {
try {
cursorEngine = new CursorSendEngine(
cursorEngine = newCursorEngine(
slotPath, actualSfMaxSegmentBytes,
actualSfMaxTotalBytes, actualSfAppendDeadlineNanos,
actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos,
actualSfSyncIntervalNanos);
} catch (SfSanitizedResidueException first) {
// NOT terminal, and it must be intercepted ahead of its
Expand All @@ -1639,9 +1633,9 @@ public Sender build() {
LOG.info("sf slot {}: sealed residue sanitized during recovery ({}); "
+ "retrying over the healed chain",
slotPath, first.getMessage());
cursorEngine = new CursorSendEngine(
cursorEngine = newCursorEngine(
slotPath, actualSfMaxSegmentBytes,
actualSfMaxTotalBytes, actualSfAppendDeadlineNanos,
actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos,
actualSfSyncIntervalNanos);
}
} catch (UnreplayableSlotException | SfRecoveryException
Expand Down Expand Up @@ -1669,7 +1663,7 @@ public Sender build() {
quarantined = true;
cursorEngine = quarantineTornSlot(
null, e, sfDir, senderId, slotPath, actualSfMaxSegmentBytes,
actualSfMaxTotalBytes, actualSfAppendDeadlineNanos,
actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos,
actualSfSyncIntervalNanos, errorHandler);
}
int actualErrorInboxCapacity = errorInboxCapacity != PARAMETER_NOT_SET_EXPLICITLY
Expand Down Expand Up @@ -1742,7 +1736,7 @@ public Sender build() {
quarantined = true;
cursorEngine = quarantineTornSlot(
cursorEngine, e, sfDir, senderId, slotPath, actualSfMaxSegmentBytes,
actualSfMaxTotalBytes, actualSfAppendDeadlineNanos,
actualSfMaxTotalBytes, sfSharedBudget, actualSfAppendDeadlineNanos,
actualSfSyncIntervalNanos, errorHandler);
} catch (Throwable t) {
// connect() failed before ownership of cursorEngine
Expand Down Expand Up @@ -2899,6 +2893,19 @@ public String getConfiguredSfDir() {
return sfDir;
}

/**
* The {@code sf_max_total_bytes} cap {@link #build()} applies: the
* configured value, or the default for the configured backing (128 MiB
* in memory, 10 GiB with {@code sf_dir}, and never less than two
* segments). Introspection hook for the connection pool, which sizes
* the one budget all of its senders share from this value. Returns
* {@code -1} for transports other than WebSocket, which have no
* segment ring to cap.
*/
public long getResolvedSfMaxTotalBytes() {
return protocol == PROTOCOL_WEBSOCKET ? resolveSfMaxTotalBytes() : -1L;
}

/**
* Excludes the connection pool's <em>live</em> slot set from
* {@link #drainOrphans(boolean)} scanning: a sibling slot under
Expand Down Expand Up @@ -3035,13 +3042,19 @@ public LineSenderBuilder storeAndForwardMaxSegmentBytes(long maxSegmentBytes) {
}

/**
* Hard cap on cursor-allocated bytes (active + spare + sealed
* segments). When the cap is reached, the producer's
* {@code Sender.flush()} blocks until ACK-driven trim frees space;
* if the cap is exhausted past the configured deadline (default 30 s),
* {@code flush()} throws. Default: {@code 128 MiB}, which applies to
* both memory-mode and SF-mode rings — for SF deployments with
* cheap disk, raise this knob explicitly. WebSocket transport only.
* Hard cap on cursor-allocated bytes: active + spare + sealed
* segments, plus the symbol-dictionary side-file with {@code sf_dir}.
* When the cap is reached, the producer's {@code Sender.flush()}
* blocks until ACK-driven trim frees space; if the cap is exhausted
* past the configured deadline (default 30 s), {@code flush()} throws.
* Default: {@code 128 MiB} of native memory without {@code sf_dir},
* {@code 10 GiB} of disk with it. WebSocket transport only.
* <p>
* The cap belongs to this sender alone unless it is built with
* {@link #storeAndForwardSharedBudget(SegmentBudget)}. The connection
* pool behind {@link QuestDB} does that for every sender it builds,
* so in a pool the cap bounds all pooled senders together, not each
* of them.
*/
public LineSenderBuilder storeAndForwardMaxTotalBytes(long maxTotalBytes) {
if (protocol != PARAMETER_NOT_SET_EXPLICITLY && protocol != PROTOCOL_WEBSOCKET) {
Expand All @@ -3054,6 +3067,29 @@ public LineSenderBuilder storeAndForwardMaxTotalBytes(long maxTotalBytes) {
return this;
}

/**
* Charges this sender's cursor segments to {@code budget}, shared with
* every other sender built on it, instead of to a private budget of
* {@code sf_max_total_bytes}. Together those senders never hold more
* than the budget's capacity, except that each is always granted its
* minimum working set: the active segment plus one spare. This
* sender's own {@code sf_max_total_bytes} then caps only the
* background orphan drainers it starts. The connection pool behind
* {@link QuestDB} builds every sender this way. WebSocket transport
* only.
*
* @param budget the shared budget, or {@code null} for a private one
* (the default)
* @return this instance for method chaining
*/
public LineSenderBuilder storeAndForwardSharedBudget(SegmentBudget budget) {
if (protocol != PARAMETER_NOT_SET_EXPLICITLY && protocol != PROTOCOL_WEBSOCKET) {
throw new LineSenderException("store_and_forward is only supported for WebSocket transport");
}
this.sfSharedBudget = budget;
return this;
}

/**
* Sets the target background checkpoint cadence for
* {@link SfDurability#PERIODIC}. Scheduler and storage latency add to
Expand Down Expand Up @@ -3090,6 +3126,25 @@ private static int getValue(CharSequence configurationString, int pos, StringSin
return pos;
}

// Every cursor engine build() creates -- including the fresh slot that
// replaces a quarantined one -- charges the shared budget when the sender
// was given one, and a private budget of sfMaxTotalBytes otherwise.
private static CursorSendEngine newCursorEngine(
String slotPath,
long sfMaxSegmentBytes,
long sfMaxTotalBytes,
SegmentBudget sfSharedBudget,
long sfAppendDeadlineNanos,
long sfSyncIntervalNanos
) {
if (sfSharedBudget != null) {
return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfSharedBudget,
sfAppendDeadlineNanos, sfSyncIntervalNanos);
}
return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes,
sfAppendDeadlineNanos, sfSyncIntervalNanos);
}

private static SfDurability parseDurabilityValue(@NotNull StringSink value) {
if (Chars.equalsIgnoreCase("memory", value)) return SfDurability.MEMORY;
if (Chars.equalsIgnoreCase("periodic", value)) return SfDurability.PERIODIC;
Expand Down Expand Up @@ -3246,8 +3301,8 @@ private static long parseSizeValue(@NotNull StringSink value, @NotNull String na
private static CursorSendEngine quarantineTornSlot(
CursorSendEngine torn, RuntimeException cause, String sfDir,
String senderId, String slotPath,
long sfMaxSegmentBytes, long sfMaxTotalBytes, long sfAppendDeadlineNanos,
long sfSyncIntervalNanos,
long sfMaxSegmentBytes, long sfMaxTotalBytes, SegmentBudget sfSharedBudget,
long sfAppendDeadlineNanos, long sfSyncIntervalNanos,
io.questdb.client.SenderErrorHandler errorHandler
) {
// The verdict, and the reason, come from the recovery seed -- the only code that
Expand Down Expand Up @@ -3322,7 +3377,7 @@ private static CursorSendEngine quarantineTornSlot(
String.valueOf(handlerFailure));
}
}
return new CursorSendEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes,
return newCursorEngine(slotPath, sfMaxSegmentBytes, sfMaxTotalBytes, sfSharedBudget,
sfAppendDeadlineNanos, sfSyncIntervalNanos);
}

Expand Down Expand Up @@ -4351,6 +4406,21 @@ private void http() {
protocol = PROTOCOL_HTTP;
}

private long resolveSfMaxSegmentBytes() {
return sfMaxSegmentBytes == PARAMETER_NOT_SET_EXPLICITLY ? DEFAULT_SEGMENT_BYTES : sfMaxSegmentBytes;
}

// The default cap depends on the backing: RAM (memory mode) is tight by
// default; disk (SF mode) is cheap, so its default is generous enough
// that normal traffic never hits it.
private long resolveSfMaxTotalBytes() {
if (sfMaxTotalBytes != PARAMETER_NOT_SET_EXPLICITLY) {
return sfMaxTotalBytes;
}
long defaultMaxTotal = sfDir == null ? DEFAULT_MAX_BYTES_MEMORY : DEFAULT_MAX_BYTES_SF;
return Math.max(defaultMaxTotal, resolveSfMaxSegmentBytes() * 2);
}

private void tcp() {
if (protocol != PARAMETER_NOT_SET_EXPLICITLY) {
throw new LineSenderException("protocol was already configured ")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ public CursorSendEngine(String sfDir, long segmentSizeBytes,
@TestOnly
public CursorSendEngine(String sfDir, long segmentSizeBytes,
long maxTotalBytes, long appendDeadlineNanos, FilesFacade dictFf) {
this(sfDir, segmentSizeBytes, null, true, maxTotalBytes,
this(sfDir, segmentSizeBytes, null, true, new SegmentBudget(maxTotalBytes),
appendDeadlineNanos, 0L, dictFf);
}

Expand All @@ -292,8 +292,23 @@ public CursorSendEngine(String sfDir, long segmentSizeBytes,
public CursorSendEngine(String sfDir, long segmentSizeBytes,
long maxTotalBytes, long appendDeadlineNanos,
long syncIntervalNanos) {
this(sfDir, segmentSizeBytes, null, true, maxTotalBytes,
appendDeadlineNanos, syncIntervalNanos);
this(sfDir, segmentSizeBytes, null, true, new SegmentBudget(maxTotalBytes),
appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE);
}

/**
* As {@link #CursorSendEngine(String, long, long, long, long)}, but the
* engine's private {@link SegmentManager} charges its segments to
* {@code budget} instead of to a budget of its own. Engines built on one
* budget share its capacity: together they never hold more than it, apart
* from the minimum working set (the active segment plus one spare) each
* ring is always granted. A sender pool builds every sender this way, so
* {@code sf_max_total_bytes} caps the pool rather than each pooled sender.
*/
public CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentBudget budget,
long appendDeadlineNanos, long syncIntervalNanos) {
this(sfDir, segmentSizeBytes, null, true, budget,
appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE);
}

/**
Expand All @@ -319,19 +334,16 @@ private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager man
private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager,
boolean ownsManager, long appendDeadlineNanos,
long syncIntervalNanos) {
this(sfDir, segmentSizeBytes, manager, ownsManager,
SegmentManager.UNLIMITED_TOTAL_BYTES, appendDeadlineNanos, syncIntervalNanos);
}

private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager,
boolean ownsManager, long maxTotalBytes, long appendDeadlineNanos,
long syncIntervalNanos) {
this(sfDir, segmentSizeBytes, manager, ownsManager, maxTotalBytes,
// A caller-supplied manager already charges its own budget.
this(sfDir, segmentSizeBytes, manager, ownsManager, null,
appendDeadlineNanos, syncIntervalNanos, FilesFacade.INSTANCE);
}

// budget is what an owned manager (ownsManager && manager == null) charges
// its segments to; unused, and may be null, when the caller supplies the
// manager.
private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager manager,
boolean ownsManager, long maxTotalBytes, long appendDeadlineNanos,
boolean ownsManager, SegmentBudget budget, long appendDeadlineNanos,
long syncIntervalNanos, FilesFacade dictFf) {
this.dictFf = dictFf;
// Allocate the bound callback before constructing an owned manager.
Expand All @@ -355,7 +367,7 @@ private CursorSendEngine(String sfDir, long segmentSizeBytes, SegmentManager man
}
if (ownsManager && manager == null) {
manager = new SegmentManager(
segmentSizeBytes, SegmentManager.DEFAULT_POLL_NANOS, maxTotalBytes);
segmentSizeBytes, SegmentManager.DEFAULT_POLL_NANOS, budget);
}
SlotLock acquiredLock = null;
if (!memoryMode) {
Expand Down
Loading
Loading