{"record":{"id":"cae6b9ac71d7c75f","repo":"apache/seatunnel","slug":"receive-failed","errorCode":"RECEIVE_FAILED","errorMessage":"Interrupted while buffering MQTT source message","messagePattern":"Interrupted while buffering MQTT source message","errorType":"error_code","errorClass":"MqttConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java","lineNumber":202,"sourceCode":"    public void messageArrived(String topic, MqttMessage message) throws Exception {\n        if (message == null || message.getPayload() == null) {\n            return;\n        }\n        byte[] payload = Arrays.copyOf(message.getPayload(), message.getPayload().length);\n        try {\n            if (!messageQueue.offer(payload, QUEUE_OFFER_TIMEOUT_MS, TimeUnit.MILLISECONDS)) {\n                MqttConnectorException exception =\n                        new MqttConnectorException(\n                                MqttConnectorErrorCode.RECEIVE_FAILED,\n                                \"MQTT source message queue is full. Increase max_queue_size \"\n                                        + \"or reduce MQTT message throughput.\");\n                receiveException = exception;\n                throw exception;\n            }\n        } catch (InterruptedException e) {\n            Thread.currentThread().interrupt();\n            receiveException = e;\n            throw new MqttConnectorException(\n                    MqttConnectorErrorCode.RECEIVE_FAILED,\n                    \"Interrupted while buffering MQTT source message\",\n                    e);\n        }\n    }\n\n    @Override\n    public void deliveryComplete(IMqttDeliveryToken token) {\n        // Source-only client — outbound delivery acknowledgements are not expected.\n    }\n\n    void subscribeTopic() throws MqttException {\n        mqttClient.subscribe(sourceConfig.getTopic(), sourceConfig.getQos());\n    }\n\n    private void checkReceiveException() {\n        if (receiveException == null) {\n            return;","sourceCodeStart":184,"sourceCodeEnd":220,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-mqtt/src/main/java/org/apache/seatunnel/connectors/seatunnel/mqtt/source/MqttSourceReader.java#L184-L220","documentation":"In messageArrived (the Paho callback), the reader puts incoming messages into a bounded queue. If the thread waiting to enqueue is interrupted, it restores the interrupt flag, records the exception, and throws MqttConnectorException with RECEIVE_FAILED. Interruption here means the task is being cancelled or shutting down.","triggerScenarios":"The thread executing messageArrived is interrupted while blocked on queue.put (queue full) or queue.offer — i.e. task cancellation/stop arrives while the buffer is full.","commonSituations":"Zeta task failover or job cancel while the sink side is slow; queue capacity too small so the reader blocks on put during interruption; backpressure from downstream pollNext stalls.","solutions":["This usually signals intentional shutdown — verify whether a cancel/failover was in progress","Increase the message queue capacity (source reader options) so messageArrived does not block on a full buffer","Slow down upstream publish rate or increase downstream throughput to avoid the full-queue block","Check the recorded receiveException surfaced later by pollNext (error 2122) for root cause"],"exampleFix":"// before\n} catch (InterruptedException e) {\n    Thread.currentThread().interrupt();\n    throw new MqttConnectorException(...RECEIVE_FAILED, \"Interrupted while buffering...\", e);\n}\n// after (no code change needed; ensure task shutdown is expected and queue sizing is adequate)\n","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n    reader.pollNext(...);\n} catch (MqttConnectorException e) {\n    if (e.getCause() instanceof InterruptedException) {\n        Thread.currentThread().interrupt(); // propagate shutdown\n        return;\n    }\n    throw e;\n}","preventionTips":["Size the message queue for peak publish rate","Ensure downstream pollNext keeps up (adequate parallelism)","Avoid interrupting reader threads outside of intentional shutdown","Always restore the interrupt flag after catching InterruptedException"],"tags":["mqtt","interrupted","thread","queue-full"],"backgroundTag":"thread-interrupted","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"}