feat: allow the Kafka record key deserializer to be configured - #148
Merged
Merged
Conversation
KafkaMessageReaderBuilder hardcoded a LongDeserializer for the record key,
matching the key written by the Apiary Hive Metastore listener. Properties
passed to withConsumerProperties() cannot change it, because collisions are
resolved in favour of the builder's own defaults:
consumerProperties.forEach((k, v) -> props.merge(k, v, (v1, v2) -> v1));
Reading a topic whose producer keys records with another type therefore fails
with "SerializationException: Size of data received by LongDeserializer is not
8". The consumer position does not advance past a record it cannot deserialize,
so the reader retries the same offset indefinitely and makes no progress.
Add withKeyDeserializer(String) so callers can supply the deserializer that
matches the producer. The default is unchanged, and consumer property
precedence is unchanged, so existing callers are unaffected.
The record key is never used to decode the event — events are read from the
record value alone — so the consumer key type becomes Object rather than Long,
which keeps the declared type honest for a reader configured with a different
key deserializer.
Co-authored-by: Claude Opus <noreply@anthropic.com>
rpoluri
approved these changes
Aug 31, 2026
3 tasks done
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
📝 Description
KafkaMessageReaderBuilderhardcodes aLongDeserializerfor the record key, matching the key written by the Apiary Hive Metastore listener. That is a reasonable default, but it currently cannot be changed — properties passed towithConsumerProperties()are merged with collisions resolved in favour of the builder's own defaults:A caller setting
key.deserializerthere is silently ignored, despite the README stating that additional consumer properties are configurable.Reading a topic whose producer keys records with another type (for example a
Stringkey) then fails with:This is worse than one rejected record. The consumer position does not advance past a record it cannot deserialize, so the reader retries the same offset indefinitely and makes no progress. Because the exception surfaces from
poll(), a caller that catches per-event exceptions and continues will spin rather than skip.Changes
KafkaMessageReaderBuilder.withKeyDeserializer(String)so callers can supply the deserializer matching their producer.LongDeserializer, and the precedence rule forwithConsumerProperties()is unchanged, so existing callers are unaffected. Flipping the merge so caller properties win would be a wider behavioural change for every existing caller; an explicit builder method is opt-in and reviewable.Objectrather thanLong. The record key is never used to decode the event — events are read from the record value alone — so this keeps the declared type honest when a different key deserializer is configured."key.deserializer"/"value.deserializer"string literals with theConsumerConfigconstants, and the deserializer class-name literals withClass#getName().buildConsumerProperties(), so the resulting configuration can be asserted without standing up a realKafkaConsumer.kafka-metastore-receiverREADME; add a8.2.5entry toCHANGELOG.md(the parent pom is already at8.2.5-SNAPSHOT).✅ Test plan
KafkaMessageReaderTest: default deserializer, the override,withConsumerPropertiesstill not winning, non-colliding consumer properties preserved, and null/empty validation.mvn testonkafka-metastore-receiver—Tests run: 18, Failures: 0, Errors: 0.mvn testonkafka-metastore-integration-tests—Tests run: 7, Failures: 0, Errors: 0, confirming the end-to-end listener/receiver round trip still passes with the default deserializer.mvn packagefrom the reactor root —BUILD SUCCESS.KafkaMessageReaderconstructor'sKafkaConsumer<Long, byte[]>signature.🔗 Related Issues
None.
🤖 Generated with opencode