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

  1. Align the kafka-clients version with the one the connector targets (check connector-kafka's declared dependency and use the same version).
  2. Disable exactly-once/transactional Kafka sink (set exactly_once=false / remove transaction settings) if you must keep the newer Kafka client.
  3. 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

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


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)