apache/seatunnel · error · PulsarConnectorException

SEND_MESSAGE_FAILED

SEND_MESSAGE_FAILED

Error message

Send message failed, please check previous error log for details.

What it means

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.

Solutions

  1. Check the job logs immediately before this error for the root Pulsar exception (the retained throwable)
  2. Verify the topic exists and auto-creation, or create it manually with correct partitioning
  3. Check Pulsar auth credentials/token expiry and broker connectivity from the worker nodes
  4. If messages are too large, reduce batch/message size or raise broker maxMessageSize

Example fix

// before: ignoring client-side setup issues surfaced late
sink { Pulsar { topic = "persistent://public/default/events" } }
// after: pre-create topic, verify auth, and size within limits
sink {
  Pulsar {
    topic = "persistent://public/default/events"
    // ensure client auth/token valid and message < broker maxMessageSize
  }
}
Defensive patterns

Strategy: retry

Validate before calling

// pre-flight: topic reachable and authorized before job start
try (PulsarAdmin admin = PulsarAdmin.builder()
        .serviceHttpUrl(adminUrl).authentication(auth).build()) {
    if (!admin.topics().getPartitionedTopicMetadata(topic).partitions ... exists) {
        throw new IllegalStateException("Topic missing: " + topic);
    }
}

Try / catch

try {
    writer.write(row);
} catch (PulsarConnectorException e) {
    if (e.getErrorCode() == SEND_MESSAGE_FAILED) {
        // inspect the root cause logged earlier; retry after fixing broker/auth/topic
    }
    throw e;
}

Prevention

When it happens

Trigger: 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.

Common situations: Topic auto-creation disabled and topic missing; Pulsar token/auth expired mid-job; message exceeding broker max payload; network partition between worker and broker.

Understand the failure class

Background: "API request failed": what wrapped HTTP errors from external APIs mean and how to find the real cause — this error's family across 29 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/0b5ca0d615465a49. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-pulsar/src/main/java/org/apache/seatunnel/connectors/seatunnel/pulsar/sink/PulsarSinkWriter.java:369

    }

    private void flushPendingMessages() throws IOException {
        for (Producer<byte[]> producer : producerMap.values()) {
            producer.flush();
        }

        while (pendingMessages.longValue() > 0) {
            checkSendException();
            for (Producer<byte[]> producer : producerMap.values()) {
                producer.flush();
            }
        }
    }

    private void checkSendException() {
        Throwable throwable = sendMessageException.get();
        if (throwable != null) {
            throw buildSendFailureException(throwable);
        }
    }

    private PulsarConnectorException buildSendFailureException(Throwable throwable) {
        return new PulsarConnectorException(
                PulsarConnectorErrorCode.SEND_MESSAGE_FAILED,
                "Send message failed, please check previous error log for details.",
                throwable);
    }

    private Throwable appendSuppressed(Throwable existingFailure, Throwable newFailure) {
        if (existingFailure == null) {
            return newFailure;
        }
        existingFailure.addSuppressed(newFailure);
        return existingFailure;
    }

View on GitHub (pinned to cf67b549a7)