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