{"record":{"id":"857bc948a0949589","repo":"apache/seatunnel","slug":"kafka-version-incompatible","errorCode":"KAFKA_VERSION_INCOMPATIBLE","errorMessage":"Incompatible KafkaProducer version","messagePattern":"Incompatible KafkaProducer version","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":174,"sourceCode":"    public boolean isTxnStarted() {\n        Object transactionManager = getTransactionManager();\n        return (boolean) ReflectionUtils.getField(transactionManager, \"transactionStarted\").get();\n    }\n\n    private static Object createProducerIdAndEpoch(long producerId, short epoch) {\n        try {\n            Field field =\n                    TransactionManager.class.getDeclaredField(PRODUCER_ID_AND_EPOCH_FIELD_NAME);\n            Class<?> clazz = field.getType();\n            Constructor<?> constructor = clazz.getDeclaredConstructor(Long.TYPE, Short.TYPE);\n            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(","sourceCodeStart":156,"sourceCodeEnd":192,"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#L156-L192","documentation":"KafkaInternalProducer.resumeTransaction uses reflection to instantiate Kafka's ProducerIdAndEpoch via its constructor. If the constructor signature changed in the bundled Kafka client version, the reflective call fails and KafkaConnectorException with code VERSION_INCOMPATIBLE ('Incompatible KafkaProducer version') is thrown.","triggerScenarios":"resumeTransaction(id, epoch) is invoked (Kafka exactly-once / transactional sink) and reflection to construct org.apache.kafka.common.ProducerIdAndEpoch fails with InvocationTargetException/NoSuchMethodException etc.","commonSituations":"Upgrading or swapping the Kafka clients dependency to a version whose internal ProducerIdAndEpoch constructor differs from what the connector was compiled against; running the connector with a shaded/relocated Kafka client on the classpath.","solutions":["Align the kafka-clients version with the one the connector targets (check connector-kafka's declared dependency and use the same version).","Disable exactly-once/transactional Kafka sink (set exactly_once=false / remove transaction settings) if you must keep the newer Kafka client.","Upgrade SeaTunnel to a release that supports your Kafka clients version."],"exampleFix":"// before (pom.xml)\n<artifactId>kafka-clients</artifactId><version>4.0.0</version>\n// after\n<artifactId>kafka-clients</artifactId><version>3.4.0</version> <!-- version the connector supports -->","handlingStrategy":"try-catch","validationCode":"// verify reflective constructor exists at startup before enabling transactions\nClass<?> c = Class.forName(\"org.apache.kafka.common.ProducerIdAndEpoch\");\nc.getConstructor(long.class, short.class); // throws NoSuchMethodException if incompatible","typeGuard":null,"tryCatchPattern":"try {\n    producer.resumeTransaction(producerId, epoch);\n} catch (KafkaConnectorException e) {\n    if (e.getSeaTunnelErrorCode().getCode().equals(KafkaConnectorErrorCode.VERSION_INCOMPATIBLE.getCode())) {\n        // abort job / alert: kafka-clients version not supported for exactly-once\n    }\n    throw e;\n}","preventionTips":["Keep kafka-clients at the version SeaTunnel connector-kafka declares.","Run a startup compatibility probe before enabling exactly-once sinks.","Avoid manually overriding kafka-clients in the engine classpath."],"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"}