{"record":{"id":"4b03feadaa5d377d","repo":"quarkusio/quarkus","slug":"no-kafka-record-metadata-found-on-incoming-message-4b03fe","errorCode":null,"errorMessage":"No Kafka record metadata found on incoming message","messagePattern":"No Kafka record metadata found on incoming message","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"extensions/smallrye-reactive-messaging-kafka/runtime/src/main/java/io/quarkus/smallrye/reactivemessaging/kafka/ReactiveExactlyOnceInvoker.java","lineNumber":26,"sourceCode":"import io.smallrye.reactive.messaging.kafka.api.IncomingKafkaRecordMetadata;\nimport io.smallrye.reactive.messaging.kafka.transactions.KafkaTransactions;\nimport io.smallrye.reactive.messaging.kafka.transactions.TransactionalEmitter;\n\npublic abstract class ReactiveExactlyOnceInvoker extends ExactlyOnceInvoker {\n\n    protected ReactiveExactlyOnceInvoker(String outgoingChannel) {\n        super(outgoingChannel);\n    }\n\n    @Override\n    protected Object invokeBeanInTransaction(KafkaTransactions<Object> tx, Object[] beanArgs,\n            IncomingKafkaRecordMetadata<?, ?> recordMeta, IncomingKafkaRecordBatchMetadata<?, ?> batchMeta) {\n        if (batchMeta != null) {\n            return tx.withTransaction(batchMeta, emitter -> handleResult(invokeBean(beanArgs), emitter));\n        } else if (recordMeta != null) {\n            return tx.withTransaction(recordMeta, emitter -> handleResult(invokeBean(beanArgs), emitter));\n        } else {\n            throw new IllegalStateException(\"No Kafka record metadata found on incoming message\");\n        }\n    }\n\n    private Uni<Void> handleResult(Object result, TransactionalEmitter<Object> emitter) {\n        if (result == null) {\n            return Uni.createFrom().voidItem();\n        }\n        if (result instanceof Uni<?> uni) {\n            return uni.onItem().invoke(item -> sendResult(item, emitter))\n                    .replaceWithVoid();\n        }\n        if (result instanceof CompletionStage<?> cs) {\n            return Uni.createFrom().completionStage(cs)\n                    .onItem().invoke(item -> sendResult(item, emitter))\n                    .replaceWithVoid();\n        }\n        if (result instanceof Multi<?> multi) {\n            return multi.onItem().invoke(emitter::send)","sourceCodeStart":8,"sourceCodeEnd":44,"githubUrl":"https://github.com/quarkusio/quarkus/blob/e1c734241f34c7919086ceb4c9262b4a58f6de44/extensions/smallrye-reactive-messaging-kafka/runtime/src/main/java/io/quarkus/smallrye/reactivemessaging/kafka/ReactiveExactlyOnceInvoker.java#L8-L44","documentation":"The reactive ExactlyOnceInvoker.invokeBeanInTransaction wraps the exactly-once processing in a Kafka transaction using either IncomingKafkaRecordBatchMetadata (batch consumption) or IncomingKafkaRecordMetadata (single record). If the incoming Message carries neither, no Kafka transaction can be opened and the method throws IllegalStateException, propagating a failure on the reactive pipeline.","triggerScenarios":"invokeBeanInTransaction is called with recordMeta == null and batchMeta == null — i.e. a Message without Kafka metadata reaches the exactly-once invoker (plain Message, emitter-produced, or metadata-stripped).","commonSituations":"Mapping IncomingKafkaRecord to a plain Message and losing metadata; feeding in-memory Emitters into exactly-once channels; test harnesses sending bare Messages.","solutions":["Preserve Kafka metadata when transforming messages (keep IncomingKafkaRecord, or copy its metadata).","Consume exactly-once channels only from the Kafka connector.","For batch consumption ensure the batch metadata is retained on the message."],"exampleFix":"// before\nreturn Message.of(transform(record.getPayload())); // drops metadata\n\n// after\nreturn record.withPayload(transform(record.getPayload())); // keeps IncomingKafkaRecordMetadata","handlingStrategy":"type-guard","validationCode":"if (msg.getMetadata(IncomingKafkaRecordMetadata.class).isEmpty()\n        && msg.getMetadata(IncomingKafkaRecordBatchMetadata.class).isEmpty()) {\n    throw new IllegalStateException(\"Missing Kafka metadata for exactly-once\");\n}","typeGuard":"boolean isKafkaMessage(Message<?> m) {\n    return m.getMetadata(IncomingKafkaRecordMetadata.class).isPresent()\n        || m.getMetadata(IncomingKafkaRecordBatchMetadata.class).isPresent();\n}","tryCatchPattern":"return uni.invokeBeanInTransaction(beanArgs)\n    .onFailure(IllegalStateException.class)\n    .recoverWithItem(t -> { log.warn(\"non-Kafka message reached exactly-once channel\"); return null; });","preventionTips":["Use record.withPayload(...) instead of Message.of(...) when transforming","Keep exactly-once channels exclusively Kafka-connector-fed","Test batch consumers retain IncomingKafkaRecordBatchMetadata"],"tags":["kafka","reactive-messaging","exactly-once","runtime"],"backgroundTag":"missing-kafka-record-metadata","analyzedSha":"e1c734241f34c7919086ceb4c9262b4a58f6de44","analyzedAt":"2026-09-05T17:01:29.979Z","contentChangedAt":"2026-09-05T17:01:29.979Z","schemaVersion":2},"datasetVersion":"2026-09-14T00:17:10.932Z"}