apache/seatunnel · error · KafkaConnectorException

KAFKA_GET_TRANSACTIONMANAGER_FAILED

KAFKA_GET_TRANSACTIONMANAGER_FAILED

Error message

Can't get transactionManager in KafkaProducer

What it means

KafkaInternalProducer.transactionManager() reflectively reads the private 'transactionManager' field of KafkaProducer. If ReflectionUtils.getField cannot locate the field (renamed, removed, or relocated in the current kafka-clients version), it throws KafkaConnectorException with code GET_TRANSACTIONMANAGER_FAILED ('Can't get transactionManager in KafkaProducer').

Solutions

  1. Pin kafka-clients to the version supported by the connector so the private 'transactionManager' field exists.
  2. Remove conflicting Kafka client jars from the plugin/engine classpath so only one kafka-clients version is loaded.
  3. Disable Kafka transactions (exactly_once=false) to avoid the reflective transactionManager access path.

Example fix

// before
Runtime classpath has kafka-clients 2.x and 4.x mixed; field lookup fails
// after
Exclude transitive kafka-clients and depend solely on the connector-supported version
Defensive patterns

Strategy: try-catch

Validate before calling

// early check that the reflective field exists
java.lang.reflect.Field f = org.apache.kafka.clients.producer.KafkaProducer.class.getDeclaredField("transactionManager");

Try / catch

try {
    producer.beginTransaction(...);
} catch (KafkaConnectorException e) {
    if (e.getSeaTunnelErrorCode().getCode().equals(KafkaConnectorErrorCode.GET_TRANSACTIONMANAGER_FAILED.getCode())) {
        // switch to non-transactional sink or fix kafka-clients version
    }
    throw e;
}

Prevention

When it happens

Trigger: Any transactional path (resumeTransaction, beginTransaction, etc.) that calls transactionManager() while the KafkaProducer class on the classpath has no accessible field literally named 'transactionManager'.

Common situations: Kafka clients upgrade where internals were refactored; classpath shading/relocation hiding the field; mixing Kafka client versions between the engine and the connector plugin.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/adceccaf9e8405ab. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaInternalProducer.java:185

            constructor.setAccessible(true);
            return constructor.newInstance(producerId, epoch);
        } catch (InvocationTargetException
                | InstantiationException
                | IllegalAccessException
                | NoSuchFieldException
                | NoSuchMethodException e) {
            throw new KafkaConnectorException(
                    KafkaConnectorErrorCode.VERSION_INCOMPATIBLE,
                    "Incompatible KafkaProducer version",
                    e);
        }
    }

    private Object getTransactionManager() {
        Optional<Object> transactionManagerOptional =
                ReflectionUtils.getField(this, KafkaProducer.class, "transactionManager");
        if (!transactionManagerOptional.isPresent()) {
            throw new KafkaConnectorException(
                    KafkaConnectorErrorCode.GET_TRANSACTIONMANAGER_FAILED,
                    "Can't get transactionManager in KafkaProducer");
        }
        return transactionManagerOptional.get();
    }

    private static void transitionTransactionManagerStateTo(
            Object transactionManager, String state) {
        ReflectionUtils.invoke(
                transactionManager, "transitionTo", getTransactionManagerState(state));
    }

    @SuppressWarnings({"unchecked", "rawtypes"})
    private static Enum<?> getTransactionManagerState(String enumName) {
        try {
            Class<Enum> cl = (Class<Enum>) Class.forName(TRANSACTION_MANAGER_STATE_ENUM);
            return Enum.valueOf(cl, enumName);
        } catch (ClassNotFoundException e) {

View on GitHub (pinned to cf67b549a7)