{"record":{"id":"9c64a887d847ce24","repo":"quarkusio/quarkus","slug":"no-kafka-record-metadata-found-on-incoming-message","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/ExactlyOnceInvoker.java","lineNumber":53,"sourceCode":"        Object[] beanArgs = Arrays.copyOf(args, args.length - SYNTHETIC_PARAM_COUNT);\n\n        return invokeBeanInTransaction(kafkaTransactions, beanArgs, recordMeta, batchMeta);\n    }\n\n    protected Object invokeBeanInTransaction(KafkaTransactions<Object> tx, Object[] beanArgs,\n            IncomingKafkaRecordMetadata<?, ?> recordMeta, IncomingKafkaRecordBatchMetadata<?, ?> batchMeta) {\n        if (batchMeta != null) {\n            tx.withTransactionAndAwait(batchMeta, emitter -> {\n                sendResult(invokeBean(beanArgs), emitter);\n                return Uni.createFrom().voidItem();\n            });\n        } else if (recordMeta != null) {\n            tx.withTransactionAndAwait(recordMeta, emitter -> {\n                sendResult(invokeBean(beanArgs), emitter);\n                return Uni.createFrom().voidItem();\n            });\n        } else {\n            throw new IllegalStateException(\"No Kafka record metadata found on incoming message\");\n        }\n        return null;\n    }\n\n    protected void sendResult(Object result, TransactionalEmitter<Object> emitter) {\n        if (result == null) {\n            return;\n        }\n        if (result instanceof Iterable<?> items) {\n            for (Object item : items) {\n                emitter.send(item);\n            }\n        } else if (result instanceof Multi<?> multi) {\n            for (Object item : multi.subscribe().asIterable()) {\n                emitter.send(item);\n            }\n        } else {\n            emitter.send(result);","sourceCodeStart":35,"sourceCodeEnd":71,"githubUrl":"https://github.com/quarkusio/quarkus/blob/e1c734241f34c7919086ceb4c9262b4a58f6de44/extensions/smallrye-reactive-messaging-kafka/runtime/src/main/java/io/quarkus/smallrye/reactivemessaging/kafka/ExactlyOnceInvoker.java#L35-L71","documentation":"The blocking ExactlyOnceInvoker sends the produced result inside a Kafka transaction that must be driven by the record metadata of the incoming message (IncomingKafkaRecordMetadata). When neither record nor batch metadata is present on the incoming message, no Kafka transaction can be created, so the invoker throws IllegalStateException at runtime.","triggerScenarios":"invoke() reaches invokeBeanInTransaction with a Message lacking IncomingKafkaRecordMetadata and IncomingKafkaRecordBatchMetadata — typically a plain Message created by application code or bridged from another connector into an exactly-once channel.","commonSituations":"Emitting plain Message payloads via an Emitter into a channel consumed by an @ExactlyOnce method, or transforming messages and dropping Kafka metadata (e.g. Message.of(payload) mapping).","solutions":["Ensure messages arriving at the exactly-once method carry IncomingKafkaRecordMetadata (consume directly from the Kafka connector, don't rewrite the Message).","When transforming, use IncomingKafkaRecord/Message.withMetadataWithFallback to preserve metadata.","Do not feed Emitters/in-memory channels into @ExactlyOnce consumers."],"exampleFix":"// before\nMessage<String> m = Message.of(payload); // metadata lost\n\n// after\nMessage<String> m = message.withMetadataWithFallback(original.getMetadata()); // preserve Kafka metadata","handlingStrategy":"try-catch","validationCode":"// guard before invoking exactly-once logic\nIncomingKafkaRecordMetadata<?, ?> meta = msg.getMetadata(IncomingKafkaRecordMetadata.class).orElse(null);\nif (meta == null && msg.getMetadata(IncomingKafkaRecordBatchMetadata.class).isEmpty()) {\n    throw new IllegalStateException(\"Message lacks Kafka metadata; cannot do exactly-once\");\n}","typeGuard":"Optional<IncomingKafkaRecordMetadata<?, ?>> hasKafkaMeta(Message<?> m) {\n    return m.getMetadata(IncomingKafkaRecordMetadata.class);\n}","tryCatchPattern":"try {\n    invoker.invoke(message);\n} catch (IllegalStateException e) {\n    if (e.getMessage().contains(\"No Kafka record metadata\")) {\n        // nack or route to dead-letter; exactly-once cannot proceed\n        message.nack(e);\n    } else throw e;\n}","preventionTips":["Never rewrite Messages with Message.of() when transforming Kafka records; preserve metadata","Do not feed Emitters/in-memory channels into @ExactlyOnce consumers","Add a unit test asserting IncomingKafkaRecordMetadata presence on messages"],"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"}