{"record":{"id":"3cc3b3841832b524","repo":"apache/seatunnel","slug":"pulsarconnectorerrorcode-create-transaction-failed","errorCode":"PulsarConnectorErrorCode.CREATE_TRANSACTION_FAILED","errorMessage":"Pulsar transaction create fail.","messagePattern":"Pulsar transaction create fail\\.","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":346,"sourceCode":"                if (!rowTypeFieldNames.contains(partitionKeyField)) {\n                    throw new PulsarConnectorException(\n                            CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,\n                            String.format(\n                                    \"Partition key field not found: %s, rowType: %s\",\n                                    partitionKeyField, rowTypeFieldNames));\n                }\n            }\n            return partitionKeyFields;\n        }\n        return Collections.emptyList();\n    }\n\n    private TransactionImpl createTransaction() {\n        try {\n            return (TransactionImpl)\n                    PulsarConfigUtil.getTransaction(pulsarClient, transactionTimeout);\n        } catch (Exception e) {\n            throw new PulsarConnectorException(\n                    PulsarConnectorErrorCode.CREATE_TRANSACTION_FAILED,\n                    \"Pulsar transaction create fail.\",\n                    e);\n        }\n    }\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    }","sourceCodeStart":328,"sourceCodeEnd":364,"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#L328-L364","documentation":"The writer creates a Pulsar transaction via PulsarConfigUtil.getTransaction with the configured transaction timeout. If the client or broker fails to open the transaction (transactions not enabled on the broker, connectivity, timeout), the exception is wrapped in CREATE_TRANSACTION_FAILED.","triggerScenarios":"Exactly-once sink path calls createTransaction (constructor or snapshotState); PulsarConfigUtil.getTransaction throws because broker transactions are disabled, the client can't reach the broker, or the timeout is invalid.","commonSituations":"Broker started without transactionCoordinatorEnabled=true; namespace-level transaction TTL/timeout too low; network partition; requesting transactions with a Pulsar version that doesn't support them.","solutions":["Enable transactions on the broker (transactionCoordinatorEnabled=true in broker.conf) and ensure cluster supports them.","Increase/validate the transaction timeout option to a sane value (within broker's max allowed).","Check the wrapped cause: connectivity errors need broker/network fixes; unsupported-operation errors need a Pulsar upgrade or disabling exactly-once (use at-least-once)."],"exampleFix":"// before\nsink {\n  Pulsar {\n    transaction_timeout = 1\n  }\n}\n// after\nsink {\n  Pulsar {\n    transaction_timeout = 3600000\n  }\n}","handlingStrategy":"try-catch","validationCode":"// verify broker transaction support before enabling exactly-once\n// broker.conf: transactionCoordinatorEnabled=true","typeGuard":null,"tryCatchPattern":"try {\n    TransactionImpl txn = createTransaction();\n} catch (PulsarConnectorException e) {\n    log.error(\"Transaction creation failed; check broker transactionCoordinatorEnabled and timeout\", e);\n    throw e;\n}","preventionTips":["Enable transaction coordinator on the Pulsar broker.","Set a valid transaction timeout within broker limits.","Fall back to at-least-once semantics if transactions are unavailable."],"tags":["pulsar","transaction","broker-config"],"backgroundTag":"invalid-config-value","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"}