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
- Pin kafka-clients to the version supported by the connector so the private 'transactionManager' field exists.
- Remove conflicting Kafka client jars from the plugin/engine classpath so only one kafka-clients version is loaded.
- 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
- Do not shade/relocate kafka-clients for jobs using transactional Kafka sinks.
- Ensure exactly one kafka-clients version on the plugin classpath.
- Pin dependency versions; avoid BOM-driven silent upgrades of kafka-clients.
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
- KAFKA_VERSION_INCOMPATIBLE
- KAFKA_TRANSACTION_NOT_STARTED
- Suppressed an additional asynchronous send failure of Kafka…
- COMMON-02
- COMMON-02
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)