apache/seatunnel · error · PulsarConnectorException

PulsarConnectorErrorCode.CREATE_PRODUCER_FAILED

PulsarConnectorErrorCode.CREATE_PRODUCER_FAILED

Error message

Failed to create Pulsar producer for topic: %s

What it means

The writer caches one Pulsar Producer<byte[]> per topic, created lazily by producerCreator.create(topic). If the Pulsar client fails to create the producer (broker unreachable, topic doesn't exist and auto-creation is disabled, auth failure), the underlying PulsarClientException is wrapped in this connector error with CREATE_PRODUCER_FAILED.

Solutions

  1. Verify serviceUrl and network connectivity from the worker to the Pulsar broker (and brokerClientAuthentication parameters).
  2. Pre-create the topic or enable broker auto-creation (allowAutoTopicCreation=true) or set the sink option to permit creation.
  3. Check the wrapped PulsarClientException cause for the exact broker rejection (authorization vs not-found vs timeout) and fix accordingly.

Example fix

// before
sink {
  Pulsar {
    serviceUrl = "pulsar://wrong-host:6650"
    topic = "my-topic"
  }
}
// after
sink {
  Pulsar {
    serviceUrl = "pulsar://pulsar-broker:6650"
    topic = "persistent://public/default/my-topic"
  }
}
Defensive patterns

Strategy: try-catch

Validate before calling

// pre-check broker reachability before job submission
PulsarClient ping = PulsarClient.builder().serviceUrl(serviceUrl).build();
ping.getPartitionsForTopic(topic).get();

Try / catch

try {
    producer = writer.getOrCreateProducer(topic);
} catch (PulsarConnectorException e) {
    Throwable cause = e.getCause();
    if (cause instanceof PulsarClientException.NotFoundException) { /* create topic */ }
    else if (cause instanceof PulsarClientException.AuthorizationException) { /* fix ACLs */ }
    throw e;
}

Prevention

When it happens

Trigger: First write to a topic triggers getOrCreateProducer; the Pulsar client throws (connection refused, topic not found when allowAutoTopicCreation=false, invalid topic name, authentication/authorization failure).

Common situations: Broker down or wrong serviceUrl; topic doesn't exist and broker has auto-creation disabled; namespace policy blocks creation; tenant/namespace authorization denied; DNS/network issues from the worker node.

Understand the failure class

Background: ECONNREFUSED and "connection refused" / "could not connect to server" errors: what they mean and how to fix them — this error's family across 44 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/3296fecf3f3493ec. 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:169

        }

        return topic;
    }

    Producer<byte[]> getOrCreateProducer(String topic) {
        Producer<byte[]> existing = producerMap.get(topic);
        if (existing != null) {
            return existing;
        }

        try {
            Producer<byte[]> producer = producerCreator.create(topic);

            producerMap.put(topic, producer);
            return producer;

        } catch (PulsarClientException e) {
            throw new PulsarConnectorException(
                    PulsarConnectorErrorCode.CREATE_PRODUCER_FAILED,
                    "Failed to create Pulsar producer for topic: " + topic,
                    e);
        }
    }

    @Override
    public void write(SeaTunnelRow element) throws IOException {
        checkSendException();

        String topic = resolveTopic(element);
        byte[] message = serializationSchema.serialize(element);
        byte[] key = null;
        if (keySerializationSchema != null) {
            key = keySerializationSchema.serialize(element);
        }

        Producer<byte[]> topicProducer = getOrCreateProducer(topic);

View on GitHub (pinned to cf67b549a7)