apache/seatunnel · warning

Ingress queue physically full, returning RETRY

Error message

Ingress queue physically full, returning RETRY (capacity={}, rejectCount={}, batchId={})

What it means

WARN log emitted by handleBatchRecord when the ingress queue is not merely past the watermark but physically full (enqueue failed), so the batch is rejected outright with the RETRY response code. Unlike the watermark path, the queue cannot accept any more records until the reader drains it. The client must resend the batch.

Solutions

  1. Have the edge client honor the RETRY response with exponential backoff instead of tight-loop resending
  2. Increase local-queue-capacity
  3. Investigate why the reader is not draining: check sink throughput, checkpoint/flush cadence, thread CPU
  4. Scale the job (more source parallelism or engine resources) to match ingress rate
  5. Rate-limit or batch on the device side so average send rate <= consumption rate

Example fix

// before: client immediately resends on RETRY
while (!send(batch)) { /* tight loop */ }
// after: backoff before resending
long backoffMs = Math.min(60000, base * (1L << attempts));
Thread.sleep(backoffMs);
send(batch);
Defensive patterns

Strategy: retry

Validate before calling

// exponential backoff before resend on RETRY
long delay = Math.min(60_000, 500L * (1L << attempt));
Thread.sleep(delay);

Prevention

When it happens

Trigger: handleBatchRecord called when the internal queue's remaining capacity is zero (enqueue fails), after incrementing queueFullCount; logged on first and every 100th consecutive rejection.

Common situations: Sustained producer rate exceeding reader consumption for a long period; reader thread blocked or slow (downstream congestion); queue capacity too small for steady-state throughput; edge devices retrying immediately without honoring backoff, keeping the queue pinned full.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/a7c9152257fd54e4. 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:193

                }
                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);
                    return EdgeSocketResponseCode.RECEIVED.getCode();
                }
            }
            queueFullCount++;
            if (queueFullCount == 1 || queueFullCount % 100 == 0) {
                log.warn(
                        "Ingress queue physically full, returning RETRY "
                                + "(capacity={}, rejectCount={}, batchId={})",
                        config.getLocalQueueCapacity(),
                        queueFullCount,
                        batchId);
            }
            return EdgeSocketResponseCode.RETRY.getCode();
        } catch (EdgeSocketConnectorException connectorException) {
            if (isDecryptionError(connectorException)) {
                log.warn(
                        "Decryption failed for batchId={}, check secret_key configuration",
                        batchId,
                        connectorException);
                return EdgeSocketResponseCode.DECRYPT_FAILED.getCode();
            }
            log.warn("Decode ingress packet failed for batchId={}", batchId, connectorException);
            return EdgeSocketResponseCode.DECODE_FAILED.getCode();
        } catch (Exception decodeException) {
            log.warn("Decode or enqueue ingress packet failed", decodeException);

View on GitHub (pinned to cf67b549a7)