{"record":{"id":"fba8d257fa70674f","repo":"apache/seatunnel","slug":"ingress-queue-at-backpressure-watermark-returning","errorCode":null,"errorMessage":"Ingress queue at backpressure watermark, returning QUEUE_FULL:{}ms (capacity={}, watermarkRatio={}, rejectCount={}, batchId={})","messagePattern":"Ingress queue at backpressure watermark, returning QUEUE_FULL:(.+?)ms \\(capacity=(.+?), watermarkRatio=(.+?), rejectCount=(.+?), batchId=(.+?)\\)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketSourceReader.java","lineNumber":168,"sourceCode":"            sourceState.notifyCheckpointComplete(checkpointId);\n        }\n    }\n\n    @Override\n    public void notifyCheckpointAborted(long checkpointId) {\n        synchronized (stateLock) {\n            sourceState.notifyCheckpointAborted(checkpointId);\n        }\n    }\n\n    @Override\n    public String handleBatchRecord(long batchId, String payload) {\n        synchronized (stateLock) {\n            if (recordQueue.isBackpressure()) {\n                queueFullCount++;\n                if (queueFullCount == 1 || queueFullCount % 100 == 0) {\n                    log.warn(\n                            \"Ingress queue at backpressure watermark, returning QUEUE_FULL:{}ms \"\n                                    + \"(capacity={}, watermarkRatio={}, rejectCount={}, batchId={})\",\n                            config.getQueueFullRetryAfterMs(),\n                            config.getLocalQueueCapacity(),\n                            config.getQueueBackpressureWatermarkRatio(),\n                            queueFullCount,\n                            batchId);\n                }\n                return EdgeSocketResponseCode.QUEUE_FULL.withPayload(\n                        config.getQueueFullRetryAfterMs());\n            }\n        }\n        try {\n            EdgeSocketQueuedRecord decoded = recordDeserializer.deserializeRecord(payload);\n            decoded.setBatchId(batchId);\n            synchronized (stateLock) {\n                QueueOfferResult offerResult = recordQueue.offer(decoded);\n                if (offerResult == QueueOfferResult.ACCEPTED) {\n                    sourceState.markRecordReceived(batchId);","sourceCodeStart":150,"sourceCodeEnd":186,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketSourceReader.java#L150-L186","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Reduce the producer's send rate or add client-side batching/jitter so bursts fit within the retry-after interval","Increase local-queue-capacity in the EdgeSocket source options to absorb bursts","Raise queue-backpressure-watermark-ratio (closer to 1.0) if earlier warnings are acceptable and memory allows","Check downstream sink throughput for the real bottleneck (slow sink drains the queue slowly)","Scale source parallelism so more reader tasks drain ingress queues"],"exampleFix":"# before\nEdgeSocket {\n  local-queue-capacity = 1000\n  queue-backpressure-watermark-ratio = 0.5\n}\n# after\nEdgeSocket {\n  local-queue-capacity = 10000\n  queue-backpressure-watermark-ratio = 0.8\n}","handlingStrategy":"retry","validationCode":"// producer-side: check last QUEUE_FULL response before sending next batch\nif (\"QUEUE_FULL\".equals(lastResponseCode)) {\n    Thread.sleep(retryAfterMs);\n}","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Size local-queue-capacity for peak burst, not average rate","Alert on repeated QUEUE_FULL warnings (queueFullCount % 100 logs) as a saturation signal","Monitor queue watermark ratio in metrics","Add client-side backoff honoring retry-after"],"tags":["backpressure","queue-full","socket-source","throttling"],"backgroundTag":"rate-limit-exceeded","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}