{"record":{"id":"a7c9152257fd54e4","repo":"apache/seatunnel","slug":"ingress-queue-physically-full-returning-retry-ca","errorCode":null,"errorMessage":"Ingress queue physically full, returning RETRY (capacity={}, rejectCount={}, batchId={})","messagePattern":"Ingress queue physically full, returning RETRY \\(capacity=(.+?), 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":193,"sourceCode":"                }\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);\n                    return EdgeSocketResponseCode.RECEIVED.getCode();\n                }\n            }\n            queueFullCount++;\n            if (queueFullCount == 1 || queueFullCount % 100 == 0) {\n                log.warn(\n                        \"Ingress queue physically full, returning RETRY \"\n                                + \"(capacity={}, rejectCount={}, batchId={})\",\n                        config.getLocalQueueCapacity(),\n                        queueFullCount,\n                        batchId);\n            }\n            return EdgeSocketResponseCode.RETRY.getCode();\n        } catch (EdgeSocketConnectorException connectorException) {\n            if (isDecryptionError(connectorException)) {\n                log.warn(\n                        \"Decryption failed for batchId={}, check secret_key configuration\",\n                        batchId,\n                        connectorException);\n                return EdgeSocketResponseCode.DECRYPT_FAILED.getCode();\n            }\n            log.warn(\"Decode ingress packet failed for batchId={}\", batchId, connectorException);\n            return EdgeSocketResponseCode.DECODE_FAILED.getCode();\n        } catch (Exception decodeException) {\n            log.warn(\"Decode or enqueue ingress packet failed\", decodeException);","sourceCodeStart":175,"sourceCodeEnd":211,"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#L175-L211","documentation":"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.","triggerScenarios":"handleBatchRecord called when the internal queue's remaining capacity is zero (enqueue fails), after incrementing queueFullCount; logged on first and every 100th consecutive rejection.","commonSituations":"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.","solutions":["Have the edge client honor the RETRY response with exponential backoff instead of tight-loop resending","Increase local-queue-capacity","Investigate why the reader is not draining: check sink throughput, checkpoint/flush cadence, thread CPU","Scale the job (more source parallelism or engine resources) to match ingress rate","Rate-limit or batch on the device side so average send rate <= consumption rate"],"exampleFix":"// before: client immediately resends on RETRY\nwhile (!send(batch)) { /* tight loop */ }\n// after: backoff before resending\nlong backoffMs = Math.min(60000, base * (1L << attempts));\nThread.sleep(backoffMs);\nsend(batch);","handlingStrategy":"retry","validationCode":"// exponential backoff before resend on RETRY\nlong delay = Math.min(60_000, 500L * (1L << attempt));\nThread.sleep(delay);","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Never tight-loop resend on RETRY; always back off","Match producer rate to consumer throughput","Increase queue capacity for bursty devices","Watch for reader stalls (sink backpressure) that pin the queue full"],"tags":["backpressure","queue-full","retry","socket-source"],"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"}