Skip to content
Open
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: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,10 @@ 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.5 - 2026-08-05
### Added
- `apiary-gluesync-listener`: optional Glue client retry with exponential backoff, via `GlueClientFactory`. Off by default (safe for HMS threads); enable for Dronefly with `GLUE_RETRY_ENABLED=true` (`GLUE_RETRY_MAX_ATTEMPTS`, default `3`). Retries `ConcurrentModificationException` by default in addition to the AWS SDK's default transient error retries; configure which exception(s) to retry via `GLUE_RETRY_EXCEPTIONS` (comma-separated). Each retry is recorded via a new `glue_listener_retry_attempt` Micrometer counter tagged by `exception_type`.

## 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.
Expand Down
3 changes: 3 additions & 0 deletions hive-event-listeners/apiary-gluesync-listener/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,9 @@ The GlueSync listener can be configured by setting the following System Environm
GLUE_PREFIX|No|Prefix added to Glue databases to handle database name collisions when synchronizing multiple metastores to the Glue catalog.
ENABLE_HIVE_TO_GLUE_RENAME_OPERATION|No|Set to true in case you would like to enable Hive table renames when syncing into Glue. Default value is false.
GLUE_SKIP_ARCHIVE|No|Default value applied to `SkipArchive` on AWS Glue `UpdateTable` requests when a table does not set the `apiary.gluesync.skipArchive` property. Accepts `true` or `false` (case-insensitive). When unset, the built-in default (`true`) is used. Any other value causes the listener to fail on startup.
GLUE_RETRY_ENABLED|No|Set to `true` to also retry the exception(s) named in `GLUE_RETRY_EXCEPTIONS` (in addition to the AWS SDK's default throttle/5xx retries) with exponential backoff. Off by default so it doesn't block HMS threads; intended for the Dronefly/CLI deployment. Default value is `false`.
GLUE_RETRY_MAX_ATTEMPTS|No|Maximum number of retries per Glue API call when `GLUE_RETRY_ENABLED=true`. Default value is `3`.
GLUE_RETRY_EXCEPTIONS|No|Comma-separated list of AWS Glue exception names to retry when `GLUE_RETRY_ENABLED=true`, e.g. `ConcurrentModificationException,ThrottlingException`. Default value is `ConcurrentModificationException`.

## Table update SkipArchive
[AWS default](https://docs.aws.amazon.com/glue/latest/webapi/API_UpdateTable.html#Glue-UpdateTable-request-SkipArchive) is to archive the table on every update. With Iceberg tables this can lead to a lot of table versions. In Glue you can only have a certain limit of the number of versions and you'll get exceptions when trying to update a table once you hit that limit. Manual version removal through AWS api is then needed. To counter this the listener defaults to `skipArchive=true`, so it does *not* make an archive of the table when updating.
Expand Down
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 @@ -13,7 +13,6 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/

package com.expediagroup.apiary.extensions.gluesync.cli;

import java.util.ArrayList;
Expand All @@ -35,11 +34,12 @@

import com.amazonaws.ClientConfiguration;
import com.amazonaws.services.glue.AWSGlue;
import com.amazonaws.services.glue.AWSGlueClientBuilder;

import com.expediagroup.apiary.extensions.events.metastore.consumer.common.thrift.ThriftHiveClient;
import com.expediagroup.apiary.extensions.events.metastore.consumer.common.thrift.ThriftHiveClientFactory;
import com.expediagroup.apiary.extensions.gluesync.listener.ApiaryGlueSync;
import com.expediagroup.apiary.extensions.gluesync.listener.GlueClientFactory;
import com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricService;
import com.expediagroup.apiary.extensions.gluesync.listener.service.GlueDatabaseService;
import com.expediagroup.apiary.extensions.gluesync.listener.service.GluePartitionService;
import com.expediagroup.apiary.extensions.gluesync.listener.service.GlueTableService;
Expand All @@ -64,19 +64,17 @@ public class GlueSyncCli {
private IsIcebergTablePredicate isIcebergTablePredicate;

public GlueSyncCli() {
MetricService metricService = new MetricService();
ClientConfiguration clientConfig = new ClientConfiguration();
clientConfig.setRequestTimeout(600000);
AWSGlue glueClient = AWSGlueClientBuilder.standard()
.withRegion(System.getenv("AWS_REGION"))
.withClientConfiguration(clientConfig)
.build();
String gluePrefix = System.getenv("GLUE_PREFIX");
AWSGlue glueClient = GlueClientFactory.buildClient(System.getenv("AWS_REGION"), clientConfig, metricService);
this.thriftHiveClientFactory = new ThriftHiveClientFactory();
thriftHiveClient = thriftHiveClientFactory.newInstance(THRIFT_CONNECTION_URI, THRIFT_CONNECTION_TIMEOUT);
metastoreClient = thriftHiveClient.getMetaStoreClient();
Configuration config = new Configuration();
this.apiaryGlueSync = new ApiaryGlueSync(config, true);
this.apiaryGlueSync = new ApiaryGlueSync(config, glueClient, gluePrefix, metricService, true);
this.isIcebergTablePredicate = new IsIcebergTablePredicate();
String gluePrefix = System.getenv("GLUE_PREFIX");
this.gluePartitionService = new GluePartitionService(glueClient, gluePrefix);
this.glueDatabaseService = new GlueDatabaseService(glueClient, gluePrefix);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,6 @@
import org.slf4j.LoggerFactory;

import com.amazonaws.services.glue.AWSGlue;
import com.amazonaws.services.glue.AWSGlueClientBuilder;
import com.amazonaws.services.glue.model.AlreadyExistsException;
import com.amazonaws.services.glue.model.EntityNotFoundException;

Expand Down Expand Up @@ -83,20 +82,17 @@ public ApiaryGlueSync(Configuration config) {

public ApiaryGlueSync(Configuration config, boolean throwExceptions) {
super(config);
this.glueClient = AWSGlueClientBuilder.standard().withRegion(System.getenv("AWS_REGION")).build();
this.metricService = new MetricService();
this.glueClient = GlueClientFactory.buildClient(System.getenv("AWS_REGION"), metricService);
String gluePrefix = System.getenv("GLUE_PREFIX");
this.glueDatabaseService = new GlueDatabaseService(glueClient, gluePrefix);
this.gluePartitionService = new GluePartitionService(glueClient, gluePrefix);
this.glueTableService = new GlueTableService(glueClient, gluePartitionService, gluePrefix);
this.isIcebergPredicate = new IsIcebergTablePredicate();
this.metricService = new MetricService();
this.throwExceptions = throwExceptions;
log.debug("ApiaryGlueSync created");
}

/**
* Just for testing.
*/
public ApiaryGlueSync(Configuration config, AWSGlue glueClient, String gluePrefix, MetricService metricService,
boolean throwExceptions) {
this(config, glueClient, gluePrefix, metricService, throwExceptions, null);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
/**
* 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.
* 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.apiary.extensions.gluesync.listener;

import static com.amazonaws.retry.PredefinedRetryPolicies.DEFAULT_RETRY_CONDITION;

import java.util.Arrays;
import java.util.Collections;
import java.util.Set;
import java.util.stream.Collectors;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.amazonaws.AmazonServiceException;
import com.amazonaws.ClientConfiguration;
import com.amazonaws.retry.PredefinedRetryPolicies;
import com.amazonaws.retry.RetryPolicy;
import com.amazonaws.services.glue.AWSGlue;
import com.amazonaws.services.glue.AWSGlueClientBuilder;

import com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricService;

/**
* Builds the AWSGlue client with optional extra-exception retry policy and Micrometer metrics.
*
* Extra retries are disabled by default to avoid blocking HMS threads. Enable for Dronefly (CLI) via:
* GLUE_RETRY_ENABLED=true
* GLUE_RETRY_MAX_ATTEMPTS=3 (optional, default 3 retries = 4 total calls)
* GLUE_RETRY_EXCEPTIONS=ConcurrentModificationException,SomeOtherException (optional, comma-separated,
* defaults to ConcurrentModificationException)
*/
public class GlueClientFactory {

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

static final String ENV_RETRY_ENABLED = "GLUE_RETRY_ENABLED";
static final String ENV_RETRY_MAX_ATTEMPTS = "GLUE_RETRY_MAX_ATTEMPTS";
static final String ENV_RETRY_EXCEPTIONS = "GLUE_RETRY_EXCEPTIONS";
static final int DEFAULT_MAX_RETRIES = 3;
static final String DEFAULT_RETRY_EXCEPTION = "ConcurrentModificationException";

private GlueClientFactory() {}

public static AWSGlue buildClient(String region, MetricService metricService) {
return buildClient(region, new ClientConfiguration(), metricService);
}

public static AWSGlue buildClient(String region, ClientConfiguration config, MetricService metricService) {
if (retryEnabled()) {
int maxRetries = maxRetries();
Set<String> retryExceptions = retryExceptions();
log.info("Glue client custom retry policy active: max {} retries per call (covers throttles, 5xx, and {})",
maxRetries, retryExceptions);
config.setRetryPolicy(buildRetryPolicy(maxRetries, retryExceptions, metricService));
} else {
log.info("Glue client custom retry policy not active; SDK default retries (throttles/5xx) still apply. Set {}=true to also retry {}.",
ENV_RETRY_ENABLED, DEFAULT_RETRY_EXCEPTION);
}

AWSGlueClientBuilder builder = AWSGlueClientBuilder.standard()
.withRegion(region)
.withClientConfiguration(config);

return builder.build();
}

static boolean retryEnabled() {
return "true".equalsIgnoreCase(System.getenv(ENV_RETRY_ENABLED));
}

static int maxRetries() {
String val = System.getenv(ENV_RETRY_MAX_ATTEMPTS);
if (val != null) {
try {
int n = Integer.parseInt(val.trim());
if (n > 0) {
return n;
}
log.warn("Invalid value for {}: '{}', using default {}", ENV_RETRY_MAX_ATTEMPTS, val, DEFAULT_MAX_RETRIES);
} catch (NumberFormatException ignored) {
log.warn("Invalid value for {}: '{}', using default {}", ENV_RETRY_MAX_ATTEMPTS, val, DEFAULT_MAX_RETRIES);
}
}
return DEFAULT_MAX_RETRIES;
}

static Set<String> retryExceptions() {
String val = System.getenv(ENV_RETRY_EXCEPTIONS);
if (val == null || val.trim().isEmpty()) {
return Collections.singleton(DEFAULT_RETRY_EXCEPTION);
}
Set<String> exceptions = Arrays.stream(val.split(","))
.map(String::trim)
.filter(exceptionType -> !exceptionType.isEmpty())
.collect(Collectors.toSet());
return exceptions.isEmpty() ? Collections.singleton(DEFAULT_RETRY_EXCEPTION) : exceptions;
}

static RetryPolicy buildRetryPolicy(int maxRetries, MetricService metricService) {
return buildRetryPolicy(maxRetries, Collections.singleton(DEFAULT_RETRY_EXCEPTION), metricService);
}

static RetryPolicy buildRetryPolicy(int maxRetries, Set<String> retryExceptions, MetricService metricService) {
RetryPolicy.RetryCondition condition = (request, exception, retriesAttempted) -> {
if (DEFAULT_RETRY_CONDITION.shouldRetry(request, exception, retriesAttempted)) {
recordRetry(metricService, exceptionTag(exception));
return true;
}
if (exception instanceof AmazonServiceException) {
String errorCode = ((AmazonServiceException) exception).getErrorCode();
if (errorCode != null && retryExceptions.contains(errorCode)) {
recordRetry(metricService, errorCode);
return true;
}
}
return false;
};
return new RetryPolicy(condition, PredefinedRetryPolicies.DEFAULT_BACKOFF_STRATEGY, maxRetries, true);
}

private static String exceptionTag(Exception e) {
if (e instanceof AmazonServiceException) {
String code = ((AmazonServiceException) e).getErrorCode();
return code != null ? code : e.getClass().getSimpleName();
}
return e.getClass().getSimpleName();
}

private static void recordRetry(MetricService metricService, String exceptionType) {
if (metricService != null) {
metricService.recordGlueRetryAttempt(exceptionType);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,12 @@ public class MetricConstants {
public static final String LISTENER_EVENT = "glue_listener_event";
public static final String LISTENER_TABLE_RENAME_DURATION = "glue_listener_table_rename_duration";

public static final String GLUE_RETRY_ATTEMPT = "glue_listener_retry_attempt";

public static final String TAG_OPERATION = "operation";
public static final String TAG_RESULT = "result";
public static final String TAG_OUTCOME = "outcome";
public static final String TAG_EXCEPTION = "exception_type";

public static final String RESULT_SUCCESS = "success";
public static final String RESULT_FAILURE = "failure";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@
*/
package com.expediagroup.apiary.extensions.gluesync.listener.metrics;

import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.GLUE_RETRY_ATTEMPT;
import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.LISTENER_EVENT;
import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.TAG_EXCEPTION;
import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.TAG_OPERATION;
import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.TAG_OUTCOME;
import static com.expediagroup.apiary.extensions.gluesync.listener.metrics.MetricConstants.TAG_RESULT;
Expand Down Expand Up @@ -48,6 +50,7 @@ public class MetricService {
private final MeterRegistry registry;
private final Map<String, Counter> metrics;
private final Map<String, Counter> events = new ConcurrentHashMap<>();
private final Map<String, Counter> retryAttempts = new ConcurrentHashMap<>();

public MetricService(MeterRegistry registry) {
this.registry = registry;
Expand Down Expand Up @@ -137,4 +140,16 @@ public void recordEvent(String operation, String result, String outcome) {
log.warn("Unable to record event {} {} {}", operation, result, outcome, e);
}
}

public void recordGlueRetryAttempt(String exceptionType) {
try {
retryAttempts.computeIfAbsent(exceptionType, k ->
Counter.builder(GLUE_RETRY_ATTEMPT)
.tags(TAG_EXCEPTION, exceptionType)
.register(registry))
.increment();
} catch (Exception e) {
log.warn("Unable to record Glue retry attempt {}", exceptionType, e);
}
}
}
Loading
Loading