Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
10 changes: 10 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* 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.
* <p>
* 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) -> {
Expand Down
Original file line number Diff line number Diff line change
@@ -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());
}
}
2 changes: 1 addition & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@
<maven.compiler.plugin.version>3.13.0</maven.compiler.plugin.version>
<maven.spotbugs.plugin.version>4.9.8.2</maven.spotbugs.plugin.version>

<apiary.extensions.version>8.2.0</apiary.extensions.version>
<apiary.extensions.version>8.2.5</apiary.extensions.version>
<springframework.boot.version>3.4.13</springframework.boot.version>
<maven.shade.plugin.version>3.2.4</maven.shade.plugin.version>
<jib.maven.plugin.version>3.4.3</jib.maven.plugin.version>
Expand Down
Loading