apache/seatunnel · error · KafkaConnectorException
KAFKA_VERSION_INCOMPATIBLE
KAFKA_VERSION_INCOMPATIBLE
Error message
Incompatible KafkaProducer version
What it means
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.
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.
Example fix
// before (pom.xml) <artifactId>kafka-clients</artifactId><version>4.0.0</version> // after <artifactId>kafka-clients</artifactId><version>3.4.0</version> <!-- version the connector supports -->
Defensive patterns
Strategy: try-catch
Validate before calling
// verify reflective constructor exists at startup before enabling transactions
Class<?> c = Class.forName("org.apache.kafka.common.ProducerIdAndEpoch");
c.getConstructor(long.class, short.class); // throws NoSuchMethodException if incompatible Try / catch
try {
producer.resumeTransaction(producerId, epoch);
} catch (KafkaConnectorException e) {
if (e.getSeaTunnelErrorCode().getCode().equals(KafkaConnectorErrorCode.VERSION_INCOMPATIBLE.getCode())) {
// abort job / alert: kafka-clients version not supported for exactly-once
}
throw e;
} Prevention
- 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.
When it happens
Trigger: resumeTransaction(id, epoch) is invoked (Kafka exactly-once / transactional sink) and reflection to construct org.apache.kafka.common.ProducerIdAndEpoch fails with InvocationTargetException/NoSuchMethodException etc.
Common situations: 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.
Related errors
- KAFKA_GET_TRANSACTIONMANAGER_FAILED
- 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/857bc948a0949589.
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:174
public boolean isTxnStarted() {
Object transactionManager = getTransactionManager();
return (boolean) ReflectionUtils.getField(transactionManager, "transactionStarted").get();
}
private static Object createProducerIdAndEpoch(long producerId, short epoch) {
try {
Field field =
TransactionManager.class.getDeclaredField(PRODUCER_ID_AND_EPOCH_FIELD_NAME);
Class<?> clazz = field.getType();
Constructor<?> constructor = clazz.getDeclaredConstructor(Long.TYPE, Short.TYPE);
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(View on GitHub (pinned to cf67b549a7)