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
- Verify serviceUrl and network connectivity from the worker to the Pulsar broker (and brokerClientAuthentication parameters).
- Pre-create the topic or enable broker auto-creation (allowAutoTopicCreation=true) or set the sink option to permit creation.
- 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
- Pre-create topics or enable allowAutoTopicCreation on the broker.
- Verify serviceUrl, auth plugins, and network access from worker nodes.
- Test producer creation in a smoke job before production runs.
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
- SEND_MESSAGE_FAILED
- COMMON_WRITER_OPERATION_FAILED (CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED)
- CommonErrorCode.ILLEGAL_ARGUMENT
- CommonErrorCode.ILLEGAL_ARGUMENT
- CommonErrorCode.UNSUPPORTED_DATA_TYPE
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)