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

  1. Enable transactions on the broker (transactionCoordinatorEnabled=true in broker.conf) and ensure cluster supports them.
  2. Increase/validate the transaction timeout option to a sane value (within broker's max allowed).
  3. 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

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


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)