apache/seatunnel · error · PulsarConnectorException
PulsarConnectorErrorCode.CREATE_TRANSACTION_FAILED
PulsarConnectorErrorCode.CREATE_TRANSACTION_FAILED
Error message
Pulsar transaction create fail.
What it means
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.
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).
Example fix
// before
sink {
Pulsar {
transaction_timeout = 1
}
}
// after
sink {
Pulsar {
transaction_timeout = 3600000
}
} Defensive patterns
Strategy: try-catch
Validate before calling
// verify broker transaction support before enabling exactly-once // broker.conf: transactionCoordinatorEnabled=true
Try / catch
try {
TransactionImpl txn = createTransaction();
} catch (PulsarConnectorException e) {
log.error("Transaction creation failed; check broker transactionCoordinatorEnabled and timeout", e);
throw e;
} Prevention
- 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.
When it happens
Trigger: 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.
Common situations: 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.
Understand the failure class
Background: "Invalid value" and "allowed values are" config errors: what your library rejected and how to fix it — this error's family across 41 libraries.
Related errors
- COMMON_ILLEGAL_ARGUMENT
- CommonErrorCode.ILLEGAL_ARGUMENT
- CommonErrorCode.ILLEGAL_ARGUMENT
- CommonErrorCode.UNSUPPORTED_DATA_TYPE
- CommonErrorCode.UNSUPPORTED_DATA_TYPE
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/3cc3b3841832b524.
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:346
if (!rowTypeFieldNames.contains(partitionKeyField)) {
throw new PulsarConnectorException(
CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,
String.format(
"Partition key field not found: %s, rowType: %s",
partitionKeyField, rowTypeFieldNames));
}
}
return partitionKeyFields;
}
return Collections.emptyList();
}
private TransactionImpl createTransaction() {
try {
return (TransactionImpl)
PulsarConfigUtil.getTransaction(pulsarClient, transactionTimeout);
} catch (Exception e) {
throw new PulsarConnectorException(
PulsarConnectorErrorCode.CREATE_TRANSACTION_FAILED,
"Pulsar transaction create fail.",
e);
}
}
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();
}
}
}View on GitHub (pinned to cf67b549a7)