From ab38a05b83ca9f847aab120e5197c98204a33e4f Mon Sep 17 00:00:00 2001 From: javbeltran_expedia Date: Tue, 1 Sep 2026 11:44:25 +0200 Subject: [PATCH] feat: allow the Kafka record key deserializer to be configured Drone Fly reads record keys with a LongDeserializer, the default in the Apiary KafkaMessageReader, which matches the key written by the Apiary Hive Metastore listener. A topic populated by a different producer may key its records with another type, and reading those fails on every record with: SerializationException: Size of data received by LongDeserializer is not 8 The consumer offset does not advance past a record it cannot deserialize, so Drone Fly retries the same offset indefinitely and stops processing entirely. Resolve the key deserializer from the consumer properties and pass it to the builder explicitly, defaulting to LongDeserializer so existing deployments are unaffected. It cannot be picked up from withConsumerProperties, because those values do not override the builder's own defaults; apiary-extensions 8.2.5 adds withKeyDeserializer for this purpose. The record key is not used to process the event, so any deserializer that can read the key is safe. Co-authored-by: Claude Opus --- CHANGELOG.md | 6 +++ README.md | 10 +++++ .../dronefly/app/context/CommonBeans.java | 25 ++++++++++- .../dronefly/app/context/CommonBeansTest.java | 42 +++++++++++++++++++ pom.xml | 2 +- 5 files changed, 83 insertions(+), 2 deletions(-) create mode 100644 drone-fly-app/src/test/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeansTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index ecafb1c..91db161 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](http://keepachangelog.com/en/1.0.0/) and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.html). +## [1.0.10] - 2026-08-31 +### Added +* Support for configuring the Kafka record key deserializer through `apiary.messaging.consumer.key.deserializer` (environment variable `APIARY_MESSAGING_CONSUMER_KEY_DESERIALIZER`). Defaults to `LongDeserializer`, matching the Apiary Hive Metastore listener, so existing deployments are unchanged. Topics written by a different producer can now be consumed without a `SerializationException`. +### Changed +* Upgraded `apiary-extensions` from `8.2.0` to `8.2.5`. + ## [1.0.9] - 2026-03-25 ### Changed * Migrated project to Java 21 and Spring Boot 3.4.x. diff --git a/README.md b/README.md index 61b68c6..d86d166 100644 --- a/README.md +++ b/README.md @@ -111,6 +111,16 @@ In this case we are sending the properties to Kafka's consumer to be able to con --apiary.messaging.consumer.sasl_jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required; \ --apiary.messaging.consumer.sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler +#### Kafka record key deserializer + +Drone Fly reads record keys with a `LongDeserializer` by default, matching the key written by the Apiary Hive Metastore listener. A topic populated by a different producer may key its records with another type; reading those with the default fails on every record with `SerializationException: Size of data received by LongDeserializer is not 8`, and because the consumer offset does not advance past a record it cannot be deserialized, Drone Fly stops making progress entirely. + +Set the deserializer that matches the producer: + +- apiary.messaging.consumer.key.deserializer=org.apache.kafka.common.serialization.StringDeserializer + +The record key is not used to process the event, so any deserializer that can read the key is safe. + ## Metrics Drone Fly exposes standard [JVM and Kafka metrics](https://docs.spring.io/spring-boot/docs/current/reference/htmlsingle/#production-ready-metrics-meter) using [Prometheus on Spring Boot Actuator](https://docs.spring.io/spring-boot/docs/current/reference/html/production-ready-features.html#production-ready-metrics-export-prometheus) endpoint `/actuator/prometheus`. diff --git a/drone-fly-app/src/main/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeans.java b/drone-fly-app/src/main/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeans.java index bc0df3e..cf18d43 100644 --- a/drone-fly-app/src/main/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeans.java +++ b/drone-fly-app/src/main/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeans.java @@ -1,5 +1,5 @@ /** - * Copyright (C) 2020-2025 Expedia, Inc. + * Copyright (C) 2020-2026 Expedia, Inc. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -22,6 +22,8 @@ import org.apache.hadoop.hive.conf.HiveConf; import org.apache.hadoop.hive.metastore.MetaStoreEventListener; import org.apache.hadoop.hive.metastore.api.MetaException; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.serialization.LongDeserializer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; @@ -80,11 +82,32 @@ public MessageReaderAdapter messageReaderAdapter() { Properties consumerProperties = getConsumerProperties(); KafkaMessageReader delegate = KafkaMessageReaderBuilder. builder(bootstrapServers, topicName, instanceName). + withKeyDeserializer(keyDeserializer(consumerProperties)). withConsumerProperties(consumerProperties). build(); return new MessageReaderAdapter(delegate); } + /** + * Resolves the deserializer for the Kafka record key. + *

+ * The key type is decided by whichever producer writes the topic. The Apiary Hive Metastore + * listener writes a {@code Long}, which is the default here, but a topic populated by another + * producer may use a different type. Reading a topic with the wrong deserializer fails on every + * record, and because the consumer offset does not advance past a record it cannot deserialize, + * the service makes no progress at all. + *

+ * The value has to be read out of the consumer properties and passed to the builder explicitly: + * properties given to {@code withConsumerProperties} do not override the builder's own defaults. + * + * @param consumerProperties consumer properties bound from {@value #CONSUMER_PROPERTIES_PREFIX} + * @return the configured key deserializer, or {@link LongDeserializer} when unset + */ + static String keyDeserializer(Properties consumerProperties) { + return consumerProperties + .getProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, LongDeserializer.class.getName()); + } + private Properties getConsumerProperties() { Properties consumerProperties = new Properties(); getEnvProperties().forEach((key, value) -> { diff --git a/drone-fly-app/src/test/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeansTest.java b/drone-fly-app/src/test/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeansTest.java new file mode 100644 index 0000000..3ca9d0d --- /dev/null +++ b/drone-fly-app/src/test/java/com/expediagroup/dataplatform/dronefly/app/context/CommonBeansTest.java @@ -0,0 +1,42 @@ +/** + * Copyright (C) 2020-2026 Expedia, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.expediagroup.dataplatform.dronefly.app.context; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.util.Properties; + +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.serialization.LongDeserializer; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.junit.jupiter.api.Test; + +public class CommonBeansTest { + + @Test + public void keyDeserializerDefaultsToLong() { + assertThat(CommonBeans.keyDeserializer(new Properties())).isEqualTo(LongDeserializer.class.getName()); + } + + @Test + public void keyDeserializerUsesConfiguredValue() { + Properties consumerProperties = new Properties(); + consumerProperties + .setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); + + assertThat(CommonBeans.keyDeserializer(consumerProperties)).isEqualTo(StringDeserializer.class.getName()); + } +} diff --git a/pom.xml b/pom.xml index 83ae028..20b9479 100644 --- a/pom.xml +++ b/pom.xml @@ -36,7 +36,7 @@ 3.13.0 4.9.8.2 - 8.2.0 + 8.2.5 3.4.13 3.2.4 3.4.3