{"record":{"id":"adceccaf9e8405ab","repo":"apache/seatunnel","slug":"kafka-get-transactionmanager-failed","errorCode":"KAFKA_GET_TRANSACTIONMANAGER_FAILED","errorMessage":"Can't get transactionManager in KafkaProducer","messagePattern":"Can't get transactionManager in KafkaProducer","errorType":"error_code","errorClass":"KafkaConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaInternalProducer.java","lineNumber":185,"sourceCode":"            constructor.setAccessible(true);\n            return constructor.newInstance(producerId, epoch);\n        } catch (InvocationTargetException\n                | InstantiationException\n                | IllegalAccessException\n                | NoSuchFieldException\n                | NoSuchMethodException e) {\n            throw new KafkaConnectorException(\n                    KafkaConnectorErrorCode.VERSION_INCOMPATIBLE,\n                    \"Incompatible KafkaProducer version\",\n                    e);\n        }\n    }\n\n    private Object getTransactionManager() {\n        Optional<Object> transactionManagerOptional =\n                ReflectionUtils.getField(this, KafkaProducer.class, \"transactionManager\");\n        if (!transactionManagerOptional.isPresent()) {\n            throw new KafkaConnectorException(\n                    KafkaConnectorErrorCode.GET_TRANSACTIONMANAGER_FAILED,\n                    \"Can't get transactionManager in KafkaProducer\");\n        }\n        return transactionManagerOptional.get();\n    }\n\n    private static void transitionTransactionManagerStateTo(\n            Object transactionManager, String state) {\n        ReflectionUtils.invoke(\n                transactionManager, \"transitionTo\", getTransactionManagerState(state));\n    }\n\n    @SuppressWarnings({\"unchecked\", \"rawtypes\"})\n    private static Enum<?> getTransactionManagerState(String enumName) {\n        try {\n            Class<Enum> cl = (Class<Enum>) Class.forName(TRANSACTION_MANAGER_STATE_ENUM);\n            return Enum.valueOf(cl, enumName);\n        } catch (ClassNotFoundException e) {","sourceCodeStart":167,"sourceCodeEnd":203,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/sink/KafkaInternalProducer.java#L167-L203","documentation":"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').","triggerScenarios":"Any transactional path (resumeTransaction, beginTransaction, etc.) that calls transactionManager() while the KafkaProducer class on the classpath has no accessible field literally named 'transactionManager'.","commonSituations":"Kafka clients upgrade where internals were refactored; classpath shading/relocation hiding the field; mixing Kafka client versions between the engine and the connector plugin.","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."],"exampleFix":"// before\nRuntime classpath has kafka-clients 2.x and 4.x mixed; field lookup fails\n// after\nExclude transitive kafka-clients and depend solely on the connector-supported version","handlingStrategy":"try-catch","validationCode":"// early check that the reflective field exists\njava.lang.reflect.Field f = org.apache.kafka.clients.producer.KafkaProducer.class.getDeclaredField(\"transactionManager\");","typeGuard":null,"tryCatchPattern":"try {\n    producer.beginTransaction(...);\n} catch (KafkaConnectorException e) {\n    if (e.getSeaTunnelErrorCode().getCode().equals(KafkaConnectorErrorCode.GET_TRANSACTIONMANAGER_FAILED.getCode())) {\n        // switch to non-transactional sink or fix kafka-clients version\n    }\n    throw e;\n}","preventionTips":["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."],"tags":["kafka","reflection","version-incompatibility","transactions"],"backgroundTag":"incompatible-client-version","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"}