{"record":{"id":"3296fecf3f3493ec","repo":"apache/seatunnel","slug":"pulsarconnectorerrorcode-create-producer-failed","errorCode":"PulsarConnectorErrorCode.CREATE_PRODUCER_FAILED","errorMessage":"Failed to create Pulsar producer for topic: %s","messagePattern":"Failed to create Pulsar producer for topic: (.+?)","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":169,"sourceCode":"        }\n\n        return topic;\n    }\n\n    Producer<byte[]> getOrCreateProducer(String topic) {\n        Producer<byte[]> existing = producerMap.get(topic);\n        if (existing != null) {\n            return existing;\n        }\n\n        try {\n            Producer<byte[]> producer = producerCreator.create(topic);\n\n            producerMap.put(topic, producer);\n            return producer;\n\n        } catch (PulsarClientException e) {\n            throw new PulsarConnectorException(\n                    PulsarConnectorErrorCode.CREATE_PRODUCER_FAILED,\n                    \"Failed to create Pulsar producer for topic: \" + topic,\n                    e);\n        }\n    }\n\n    @Override\n    public void write(SeaTunnelRow element) throws IOException {\n        checkSendException();\n\n        String topic = resolveTopic(element);\n        byte[] message = serializationSchema.serialize(element);\n        byte[] key = null;\n        if (keySerializationSchema != null) {\n            key = keySerializationSchema.serialize(element);\n        }\n\n        Producer<byte[]> topicProducer = getOrCreateProducer(topic);","sourceCodeStart":151,"sourceCodeEnd":187,"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#L151-L187","documentation":"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.","triggerScenarios":"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).","commonSituations":"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.","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."],"exampleFix":"// before\nsink {\n  Pulsar {\n    serviceUrl = \"pulsar://wrong-host:6650\"\n    topic = \"my-topic\"\n  }\n}\n// after\nsink {\n  Pulsar {\n    serviceUrl = \"pulsar://pulsar-broker:6650\"\n    topic = \"persistent://public/default/my-topic\"\n  }\n}","handlingStrategy":"try-catch","validationCode":"// pre-check broker reachability before job submission\nPulsarClient ping = PulsarClient.builder().serviceUrl(serviceUrl).build();\nping.getPartitionsForTopic(topic).get();","typeGuard":null,"tryCatchPattern":"try {\n    producer = writer.getOrCreateProducer(topic);\n} catch (PulsarConnectorException e) {\n    Throwable cause = e.getCause();\n    if (cause instanceof PulsarClientException.NotFoundException) { /* create topic */ }\n    else if (cause instanceof PulsarClientException.AuthorizationException) { /* fix ACLs */ }\n    throw e;\n}","preventionTips":["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."],"tags":["pulsar","producer","broker-connectivity"],"backgroundTag":"connection-refused","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"}