{"record":{"id":"0b5ca0d615465a49","repo":"apache/seatunnel","slug":"send-message-failed-0b5ca0","errorCode":"SEND_MESSAGE_FAILED","errorMessage":"Send message failed, please check previous error log for details.","messagePattern":"Send message failed, please check previous error log for details\\.","errorType":"error_code","errorClass":"PulsarConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkWriter.java","lineNumber":369,"sourceCode":"    }\n\n    private void flushPendingMessages() throws IOException {\n        for (Producer<byte[]> producer : producerMap.values()) {\n            producer.flush();\n        }\n\n        while (pendingMessages.longValue() > 0) {\n            checkSendException();\n            for (Producer<byte[]> producer : producerMap.values()) {\n                producer.flush();\n            }\n        }\n    }\n\n    private void checkSendException() {\n        Throwable throwable = sendMessageException.get();\n        if (throwable != null) {\n            throw buildSendFailureException(throwable);\n        }\n    }\n\n    private PulsarConnectorException buildSendFailureException(Throwable throwable) {\n        return new PulsarConnectorException(\n                PulsarConnectorErrorCode.SEND_MESSAGE_FAILED,\n                \"Send message failed, please check previous error log for details.\",\n                throwable);\n    }\n\n    private Throwable appendSuppressed(Throwable existingFailure, Throwable newFailure) {\n        if (existingFailure == null) {\n            return newFailure;\n        }\n        existingFailure.addSuppressed(newFailure);\n        return existingFailure;\n    }\n","sourceCodeStart":351,"sourceCodeEnd":387,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkWriter.java#L351-L387","documentation":"Pulsar sink writes are asynchronous; send failures are captured in sendMessageException and rethrown as PulsarConnectorException (SEND_MESSAGE_FAILED) at the next checkpoint (write/prepareCommit/snapshotState/flush). The message itself is generic — the root cause is in the earlier error log/throwable.","triggerScenarios":"Producer send fails asynchronously (topic doesn't exist, authorization failure, producer closed, message too large, broker unavailable) and the next write/flush/checkpoint call surfaces it via checkSendException.","commonSituations":"Topic auto-creation disabled and topic missing; Pulsar token/auth expired mid-job; message exceeding broker max payload; network partition between worker and broker.","solutions":["Check the job logs immediately before this error for the root Pulsar exception (the retained throwable)","Verify the topic exists and auto-creation, or create it manually with correct partitioning","Check Pulsar auth credentials/token expiry and broker connectivity from the worker nodes","If messages are too large, reduce batch/message size or raise broker maxMessageSize"],"exampleFix":"// before: ignoring client-side setup issues surfaced late\nsink { Pulsar { topic = \"persistent://public/default/events\" } }\n// after: pre-create topic, verify auth, and size within limits\nsink {\n  Pulsar {\n    topic = \"persistent://public/default/events\"\n    // ensure client auth/token valid and message < broker maxMessageSize\n  }\n}","handlingStrategy":"retry","validationCode":"// pre-flight: topic reachable and authorized before job start\ntry (PulsarAdmin admin = PulsarAdmin.builder()\n        .serviceHttpUrl(adminUrl).authentication(auth).build()) {\n    if (!admin.topics().getPartitionedTopicMetadata(topic).partitions ... exists) {\n        throw new IllegalStateException(\"Topic missing: \" + topic);\n    }\n}","typeGuard":null,"tryCatchPattern":"try {\n    writer.write(row);\n} catch (PulsarConnectorException e) {\n    if (e.getErrorCode() == SEND_MESSAGE_FAILED) {\n        // inspect the root cause logged earlier; retry after fixing broker/auth/topic\n    }\n    throw e;\n}","preventionTips":["Pre-create topics (or enable auto-creation) and verify permissions","Monitor Pulsar client error logs; the surfaced message is generic — root cause is upstream","Keep message sizes under broker maxMessageSize; watch token expiry on long jobs"],"tags":["pulsar","async","producer","network"],"backgroundTag":"api-request-failed","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"}