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
4 changes: 3 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,16 @@ 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).

## 8.2.0 - TBD
## 8.2.4 - 2026-07-17
### Added
- `apiary-gluesync-listener`: per-event observability via a new `glue_listener_event` Micrometer counter, tagged with `operation` (e.g. `create_table`), `result` (`success`, `failure`, `ignored`), and `outcome` (e.g. `created`, `updated`, `deleted`, `not_found`, `renamed`, exception class name). Covers all 8 HMS event handlers. A `glue_listener_table_rename_duration` timer is also recorded on every table rename.
- `kafka-metastore-receiver`: Kafka consumer metrics via Micrometer `KafkaClientMetrics`. Bind a registry at build time with `KafkaMessageReaderBuilder.withMeterRegistry(registry)`.
### Fixed
- `apiary-gluesync-listener`: fallback Glue operations (update-after-`AlreadyExistsException`, create-after-`EntityNotFoundException`) that previously escaped their catch blocks unhandled are now consistently caught, logged, and metered by the outer exception handler.
### Changed
- Bump Micrometer from `1.9.9` to `1.14.14`.
### Reverted
- Revert `Change kafka listener hashing for ALTER_TABLE events (#133)`. `ALTER_TABLE` events now hash on the current table name again. The reverted change caused ordering discrepancies in the event flow.

## 8.1.18 - 2026-06-09
### Fixed
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,15 +43,13 @@

import com.google.common.annotations.VisibleForTesting;

import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryAlterTableEvent;
import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryListenerEvent;
import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryListenerEventFactory;
import com.expediagroup.apiary.extensions.events.metastore.io.MetaStoreEventSerDe;
import com.expediagroup.apiary.extensions.events.metastore.kafka.messaging.KafkaMessage;
import com.expediagroup.apiary.extensions.events.metastore.kafka.messaging.KafkaMessageSender;

public class KafkaMetaStoreEventListener extends MetaStoreEventListener {

private static final Logger log = LoggerFactory.getLogger(KafkaMetaStoreEventListener.class);

private final MetaStoreEventSerDe eventSerDe;
Expand All @@ -60,7 +58,7 @@ public class KafkaMetaStoreEventListener extends MetaStoreEventListener {

public KafkaMetaStoreEventListener(Configuration config) {
this(config, new ApiaryListenerEventFactory(), serDeForClassName(stringProperty(config, SERDE_CLASS)),
new KafkaMessageSender(config));
new KafkaMessageSender(config));
}

@VisibleForTesting
Expand All @@ -76,19 +74,12 @@ public KafkaMetaStoreEventListener(Configuration config) {
}

private KafkaMessage withPayload(ApiaryListenerEvent event) {
String database = event.getDatabaseName();
String table = event.getTableName();
// Ensuring ALTER_TABLE events will use old table name for partition hashing
if (event instanceof ApiaryAlterTableEvent) {
database = ((ApiaryAlterTableEvent) event).getOldTable().getDbName();
table = ((ApiaryAlterTableEvent) event).getOldTable().getTableName();
}
return KafkaMessage
.builder()
.database(database)
.table(table)
.payload(eventSerDe.marshal(event))
.build();
.builder()
.database(event.getDatabaseName())
.table(event.getTableName())
.payload(eventSerDe.marshal(event))
.build();
}

@Override
Expand Down Expand Up @@ -187,4 +178,5 @@ public void onCreateFunction(CreateFunctionEvent fnEvent) {}

@Override
public void onDropFunction(DropFunctionEvent fnEvent) {}

}
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/**
* Copyright (C) 2018-2025 Expedia, Inc.
* Copyright (C) 2018-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 Down Expand Up @@ -43,15 +43,13 @@

import com.google.common.annotations.VisibleForTesting;

import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryAlterTableEvent;
import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryListenerEvent;
import com.expediagroup.apiary.extensions.events.metastore.event.ApiaryListenerEventFactory;
import com.expediagroup.apiary.extensions.events.metastore.io.MetaStoreEventSerDe;
import com.expediagroup.apiary.extensions.events.metastore.kafka.messaging.KafkaMessage;
import com.expediagroup.apiary.extensions.events.metastore.kafka.messaging.MskMessageSender;

public class MskMetaStoreEventListener extends MetaStoreEventListener {

private static final Logger log = LoggerFactory.getLogger(MskMetaStoreEventListener.class);

private final MetaStoreEventSerDe eventSerDe;
Expand All @@ -60,7 +58,7 @@ public class MskMetaStoreEventListener extends MetaStoreEventListener {

public MskMetaStoreEventListener(Configuration config) {
this(config, new ApiaryListenerEventFactory(), serDeForClassName(stringProperty(config, SERDE_CLASS)),
new MskMessageSender(config));
new MskMessageSender(config));
}

@VisibleForTesting
Expand All @@ -76,19 +74,12 @@ public MskMetaStoreEventListener(Configuration config) {
}

private KafkaMessage withPayload(ApiaryListenerEvent event) {
String database = event.getDatabaseName();
String table = event.getTableName();
// Ensuring ALTER_TABLE events will use old table name for partition hashing
if (event instanceof ApiaryAlterTableEvent) {
database = ((ApiaryAlterTableEvent) event).getOldTable().getDbName();
table = ((ApiaryAlterTableEvent) event).getOldTable().getTableName();
}
return KafkaMessage
.builder()
.database(database)
.table(table)
.payload(eventSerDe.marshal(event))
.build();
.builder()
.database(event.getDatabaseName())
.table(event.getTableName())
.payload(eventSerDe.marshal(event))
.build();
}

@Override
Expand Down Expand Up @@ -187,4 +178,5 @@ public void onCreateFunction(CreateFunctionEvent fnEvent) {}

@Override
public void onDropFunction(DropFunctionEvent fnEvent) {}

}
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,6 @@
import static org.mockito.Mockito.when;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hive.metastore.api.Table;
import org.apache.hadoop.hive.metastore.events.AddIndexEvent;
import org.apache.hadoop.hive.metastore.events.AddPartitionEvent;
import org.apache.hadoop.hive.metastore.events.AlterIndexEvent;
Expand Down Expand Up @@ -93,12 +92,8 @@ public void onCreateTable() {
public void onAlterTable() {
AlterTableEvent event = mock(AlterTableEvent.class);
ApiaryAlterTableEvent apiaryEvent = mock(ApiaryAlterTableEvent.class);
Table oldTable = mock(Table.class);
when(apiaryEvent.getOldTable()).thenReturn(oldTable);
when(apiaryEvent.getDatabaseName()).thenReturn(DATABASE);
when(apiaryEvent.getTableName()).thenReturn(TABLE);
when(oldTable.getDbName()).thenReturn(DATABASE);
when(oldTable.getTableName()).thenReturn(TABLE);
when(apiaryListenerEventFactory.create(event)).thenReturn(apiaryEvent);
listener.onAlterTable(event);
verify(kafkaMessageSender).send(any(KafkaMessage.class));
Expand Down Expand Up @@ -221,4 +216,5 @@ public void onDropFunction() {
verify(kafkaMessageSender, never()).send(any(KafkaMessage.class));
verify(eventSerDe, never()).marshal(any(ApiaryListenerEvent.class));
}

}
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/**
* Copyright (C) 2018-2025 Expedia, Inc.
* Copyright (C) 2018-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,7 +22,6 @@
import static org.mockito.Mockito.when;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hive.metastore.api.Table;
import org.apache.hadoop.hive.metastore.events.AddIndexEvent;
import org.apache.hadoop.hive.metastore.events.AddPartitionEvent;
import org.apache.hadoop.hive.metastore.events.AlterIndexEvent;
Expand Down Expand Up @@ -93,12 +92,8 @@ public void onCreateTable() {
public void onAlterTable() {
AlterTableEvent event = mock(AlterTableEvent.class);
ApiaryAlterTableEvent apiaryEvent = mock(ApiaryAlterTableEvent.class);
Table oldTable = mock(Table.class);
when(apiaryEvent.getOldTable()).thenReturn(oldTable);
when(apiaryEvent.getDatabaseName()).thenReturn(DATABASE);
when(apiaryEvent.getTableName()).thenReturn(TABLE);
when(oldTable.getDbName()).thenReturn(DATABASE);
when(oldTable.getTableName()).thenReturn(TABLE);
when(apiaryListenerEventFactory.create(event)).thenReturn(apiaryEvent);
listener.onAlterTable(event);
verify(mskMessageSender).send(any(KafkaMessage.class));
Expand Down Expand Up @@ -221,4 +216,5 @@ public void onDropFunction() {
verify(mskMessageSender, never()).send(any(KafkaMessage.class));
verify(eventSerDe, never()).marshal(any(ApiaryListenerEvent.class));
}

}
Loading