apache/seatunnel · warning

Ingress queue at backpressure watermark, returning…

Error message

Ingress queue at backpressure watermark, returning QUEUE_FULL:{}ms (capacity={}, watermarkRatio={}, rejectCount={}, batchId={})

What it means

This is a WARN log (not an exception) emitted by EdgeSocketSourceReader.handleBatchRecord when the inbound record queue has crossed its backpressure watermark. The reader still accepts or rejects based on queue state, but it logs watermark pressure and returns a QUEUE_FULL response code to the edge client with a suggested retry-after interval. It is throttling telemetry, telling the sender the connector is saturated.

Solutions

  1. Reduce the producer's send rate or add client-side batching/jitter so bursts fit within the retry-after interval
  2. Increase local-queue-capacity in the EdgeSocket source options to absorb bursts
  3. Raise queue-backpressure-watermark-ratio (closer to 1.0) if earlier warnings are acceptable and memory allows
  4. Check downstream sink throughput for the real bottleneck (slow sink drains the queue slowly)
  5. Scale source parallelism so more reader tasks drain ingress queues

Example fix

# before
EdgeSocket {
  local-queue-capacity = 1000
  queue-backpressure-watermark-ratio = 0.5
}
# after
EdgeSocket {
  local-queue-capacity = 10000
  queue-backpressure-watermark-ratio = 0.8
}
Defensive patterns

Strategy: retry

Validate before calling

// producer-side: check last QUEUE_FULL response before sending next batch
if ("QUEUE_FULL".equals(lastResponseCode)) {
    Thread.sleep(retryAfterMs);
}

Prevention

When it happens

Trigger: handleBatchRecord is called while recordQueue.isBackpressure() is true — i.e. queue size / capacity >= config.getQueueBackpressureWatermarkRatio(). The log fires on the first rejection (queueFullCount == 1) and every 100th rejection thereafter.

Common situations: Edge devices (sensors/gateways) pushing batches faster than the SeaTunnel source can drain them; downstream sink backpressure slowing consumption; queue capacity (local-queue-capacity) configured too small for bursty traffic; watermark ratio set too aggressively low.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/fba8d257fa70674f. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketSourceReader.java:168

            sourceState.notifyCheckpointComplete(checkpointId);
        }
    }

    @Override
    public void notifyCheckpointAborted(long checkpointId) {
        synchronized (stateLock) {
            sourceState.notifyCheckpointAborted(checkpointId);
        }
    }

    @Override
    public String handleBatchRecord(long batchId, String payload) {
        synchronized (stateLock) {
            if (recordQueue.isBackpressure()) {
                queueFullCount++;
                if (queueFullCount == 1 || queueFullCount % 100 == 0) {
                    log.warn(
                            "Ingress queue at backpressure watermark, returning QUEUE_FULL:{}ms "
                                    + "(capacity={}, watermarkRatio={}, rejectCount={}, batchId={})",
                            config.getQueueFullRetryAfterMs(),
                            config.getLocalQueueCapacity(),
                            config.getQueueBackpressureWatermarkRatio(),
                            queueFullCount,
                            batchId);
                }
                return EdgeSocketResponseCode.QUEUE_FULL.withPayload(
                        config.getQueueFullRetryAfterMs());
            }
        }
        try {
            EdgeSocketQueuedRecord decoded = recordDeserializer.deserializeRecord(payload);
            decoded.setBatchId(batchId);
            synchronized (stateLock) {
                QueueOfferResult offerResult = recordQueue.offer(decoded);
                if (offerResult == QueueOfferResult.ACCEPTED) {
                    sourceState.markRecordReceived(batchId);

View on GitHub (pinned to cf67b549a7)