Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[SIP] feat(iceberg): support event listener for Iceberg REST server #5002

Draft
wants to merge 6 commits into
base: main
Choose a base branch
from
Draft
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: 4 additions & 2 deletions core/src/main/java/org/apache/gravitino/GravitinoEnv.java
Original file line number Diff line number Diff line change
Expand Up @@ -307,10 +307,12 @@ public FutureGrantManager futureGrantManager() {
return futureGrantManager;
}

public void start() {
auxServiceManager.serviceStart();
public void start(boolean isGravitinoServer) {
metricsSystem.start();
eventListenerManager.start();
if (isGravitinoServer) {
auxServiceManager.serviceStart();
}
}

/** Shutdown the Gravitino environment. */
Expand Down
43 changes: 39 additions & 4 deletions core/src/main/java/org/apache/gravitino/listener/EventBus.java
Original file line number Diff line number Diff line change
Expand Up @@ -21,29 +21,42 @@

import com.google.common.annotations.VisibleForTesting;
import java.util.List;
import java.util.stream.Collectors;
import org.apache.gravitino.listener.api.EventListenerPlugin;
import org.apache.gravitino.listener.api.EventListenerPlugin.PreEventCheckException;
import org.apache.gravitino.listener.api.event.Event;
import org.apache.gravitino.listener.api.event.PreEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

/**
* The {@code EventBus} class serves as a mechanism to dispatch events to registered listeners. It
* supports both synchronous and asynchronous listeners by categorizing them into two distinct types
* within its internal management.
*/
public class EventBus {
private static final Logger LOG = LoggerFactory.getLogger(EventBus.class);

// Holds instances of EventListenerPlugin. These instances can either be
// EventListenerPluginWrapper,
// which are meant for synchronous event listening, or AsyncQueueListener, designed for
// asynchronous event processing.
private final List<EventListenerPlugin> postEventListeners;
private final List<EventListenerPlugin> preEventListeners;

/**
* Constructs an EventBus with a predefined list of event listeners.
*
* @param postEventListeners A list of {@link EventListenerPlugin} instances that are to be
* @param eventListenerPlugins A list of {@link EventListenerPlugin} instances that are to be
* registered with this EventBus for event dispatch.
*/
public EventBus(List<EventListenerPlugin> postEventListeners) {
this.postEventListeners = postEventListeners;
public EventBus(List<EventListenerPlugin> eventListenerPlugins) {
this.postEventListeners = eventListenerPlugins;
// todo: use a more general way to create preEventListeners
this.preEventListeners =
postEventListeners.stream()
.filter(eventListenerPlugin -> !(eventListenerPlugin instanceof AsyncQueueListener))
.collect(Collectors.toList());
}

/**
Expand All @@ -53,7 +66,11 @@ public EventBus(List<EventListenerPlugin> postEventListeners) {
* @param event The event to be dispatched to all registered listeners.
*/
public void dispatchEvent(Event event) {
postEventListeners.forEach(postEventListener -> postEventListener.onPostEvent(event));
if (event instanceof PreEvent) {
dispatchPreEvent((PreEvent) event);
} else {
dispatchPostEvent(event);
}
}

/**
Expand All @@ -67,4 +84,22 @@ public void dispatchEvent(Event event) {
List<EventListenerPlugin> getPostEventListeners() {
return postEventListeners;
}

public void dispatchPreEvent(PreEvent event) {
preEventListeners.forEach(
preEventListener -> {
try {
preEventListener.onPreEvent(event);
} catch (PreEventCheckException e) {
// todo transform to ValidateException for IcebergRESTServer
throw new RuntimeException(e);
} catch (Exception e) {
LOG.warn("PreEvent handle exception,", e);
}
});
}

private void dispatchPostEvent(Event event) {
postEventListeners.forEach(postEventListener -> postEventListener.onPostEvent(event));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import java.util.Map;
import org.apache.gravitino.annotation.DeveloperApi;
import org.apache.gravitino.listener.api.event.Event;
import org.apache.gravitino.listener.api.event.PreEvent;

/**
* Defines an interface for event listeners that manage the lifecycle and state of a plugin,
Expand All @@ -33,6 +34,9 @@
*/
@DeveloperApi
public interface EventListenerPlugin {

class PreEventCheckException extends RuntimeException {}

/**
* Defines the operational modes for event processing within an event listener, catering to both
* synchronous and asynchronous processing strategies. Each mode determines how events are
Expand Down Expand Up @@ -107,6 +111,10 @@ enum Mode {
*/
void onPostEvent(Event event) throws RuntimeException;

// The related operation will fail, if throwing PreEventCheckException, will ignore other
// Exceptions.
void onPreEvent(PreEvent preEvent) throws PreEventCheckException;

/**
* Specifies the default operational mode for event processing by the plugin. The default
* implementation is synchronous, but implementers can override this to utilize asynchronous
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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 org.apache.gravitino.listener.api.event;

import org.apache.gravitino.NameIdentifier;

public abstract class PreEvent extends Event {
protected PreEvent(String user, NameIdentifier identifier) {
super(user, identifier);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.gravitino.iceberg.service.IcebergExceptionMapper;
import org.apache.gravitino.iceberg.service.IcebergObjectMapperProvider;
import org.apache.gravitino.iceberg.service.metrics.IcebergMetricsManager;
import org.apache.gravitino.listener.EventBus;
import org.apache.gravitino.metrics.MetricsSystem;
import org.apache.gravitino.metrics.source.MetricsSource;
import org.apache.gravitino.server.web.HttpServerMetricsSource;
Expand Down Expand Up @@ -66,6 +67,7 @@ private void initServer(IcebergConfig icebergConfig) {
new HttpServerMetricsSource(MetricsSource.ICEBERG_REST_SERVER_METRIC_NAME, config, server);
metricsSystem.register(httpServerMetricsSource);

EventBus eventBus = GravitinoEnv.getInstance().eventBus();
icebergCatalogWrapperManager = new IcebergCatalogWrapperManager(icebergConfig.getAllConfig());
icebergMetricsManager = new IcebergMetricsManager(icebergConfig);
config.register(
Expand All @@ -74,6 +76,7 @@ private void initServer(IcebergConfig icebergConfig) {
protected void configure() {
bind(icebergCatalogWrapperManager).to(IcebergCatalogWrapperManager.class).ranked(1);
bind(icebergMetricsManager).to(IcebergMetricsManager.class).ranked(1);
bind(eventBus).to(EventBus.class).ranked(1);
}
});

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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 org.apache.gravitino.iceberg.extension;

import java.util.Map;
import org.apache.gravitino.listener.api.EventListenerPlugin;
import org.apache.gravitino.listener.api.event.Event;
import org.apache.gravitino.listener.api.event.IcebergCreateTablePostEvent;
import org.apache.gravitino.listener.api.event.IcebergCreateTablePreEvent;
import org.apache.gravitino.listener.api.event.IcebergUpdateTablePostEvent;
import org.apache.gravitino.listener.api.event.IcebergUpdateTablePreEvent;
import org.apache.gravitino.listener.api.event.PreEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class IcebergEventLogger implements EventListenerPlugin {

private static final Logger LOG = LoggerFactory.getLogger(IcebergEventLogger.class);

@Override
public void init(Map<String, String> properties) throws RuntimeException {}

@Override
public void start() throws RuntimeException {}

@Override
public void stop() throws RuntimeException {}

@Override
public void onPostEvent(Event event) throws RuntimeException {
if (event instanceof IcebergCreateTablePostEvent) {
LOG.info(
"Create table event, request: {}",
((IcebergCreateTablePostEvent) event).createTableRequest());
} else if (event instanceof IcebergUpdateTablePostEvent) {
LOG.info(
"Update table event, request: {}",
((IcebergUpdateTablePostEvent) event).updateTableRequest());
} else {
LOG.info("Unknown event: {}", event.getClass().getSimpleName());
}
}

@Override
public void onPreEvent(PreEvent preEvent) throws RuntimeException {
if (preEvent instanceof IcebergCreateTablePreEvent) {
LOG.info(
"Create table event, request: {}",
((IcebergCreateTablePreEvent) preEvent).createTableRequest());
} else if (preEvent instanceof IcebergUpdateTablePreEvent) {
LOG.info(
"Update table event, request: {}",
((IcebergUpdateTablePreEvent) preEvent).updateTableRequest());
} else {
LOG.info("Unknown event: {}", preEvent.getClass().getSimpleName());
}
}

@Override
public Mode mode() {
return Mode.ASYNC_ISOLATED;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ private void initialize() {
}

private void start() {
gravitinoEnv.start(false);
icebergRESTService.serviceStart();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,16 +37,23 @@
import javax.ws.rs.core.Context;
import javax.ws.rs.core.MediaType;
import javax.ws.rs.core.Response;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.iceberg.service.IcebergCatalogWrapperManager;
import org.apache.gravitino.iceberg.service.IcebergObjectMapper;
import org.apache.gravitino.iceberg.service.IcebergRestUtils;
import org.apache.gravitino.iceberg.service.metrics.IcebergMetricsManager;
import org.apache.gravitino.listener.EventBus;
import org.apache.gravitino.listener.api.event.IcebergCreateTablePostEvent;
import org.apache.gravitino.listener.api.event.IcebergCreateTablePreEvent;
import org.apache.gravitino.listener.api.event.IcebergUpdateTablePostEvent;
import org.apache.gravitino.listener.api.event.IcebergUpdateTablePreEvent;
import org.apache.gravitino.metrics.MetricNames;
import org.apache.iceberg.catalog.TableIdentifier;
import org.apache.iceberg.rest.RESTUtil;
import org.apache.iceberg.rest.requests.CreateTableRequest;
import org.apache.iceberg.rest.requests.ReportMetricsRequest;
import org.apache.iceberg.rest.requests.UpdateTableRequest;
import org.apache.iceberg.rest.responses.LoadTableResponse;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -61,6 +68,7 @@ public class IcebergTableOperations {
private IcebergMetricsManager icebergMetricsManager;

private ObjectMapper icebergObjectMapper;
private EventBus eventBus;

@SuppressWarnings("UnusedVariable")
@Context
Expand All @@ -69,10 +77,12 @@ public class IcebergTableOperations {
@Inject
public IcebergTableOperations(
IcebergCatalogWrapperManager icebergCatalogWrapperManager,
IcebergMetricsManager icebergMetricsManager) {
IcebergMetricsManager icebergMetricsManager,
EventBus eventBus) {
this.icebergCatalogWrapperManager = icebergCatalogWrapperManager;
this.icebergObjectMapper = IcebergObjectMapper.getInstance();
this.icebergMetricsManager = icebergMetricsManager;
this.eventBus = eventBus;
}

@GET
Expand All @@ -97,10 +107,24 @@ public Response createTable(
"Create Iceberg table, namespace: {}, create table request: {}",
namespace,
createTableRequest);
return IcebergRestUtils.ok(

eventBus.dispatchPreEvent(
new IcebergCreateTablePreEvent(
"user", NameIdentifier.of(namespace, createTableRequest.name()), createTableRequest));

LoadTableResponse loadTableResponse =
icebergCatalogWrapperManager
.getOps(prefix)
.createTable(RESTUtil.decodeNamespace(namespace), createTableRequest));
.createTable(RESTUtil.decodeNamespace(namespace), createTableRequest);

eventBus.dispatchEvent(
Copy link
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should be refactored to eventBus.dispatchPostEvent

new IcebergCreateTablePostEvent(
"user",
NameIdentifier.of(namespace, createTableRequest.name()),
createTableRequest,
loadTableResponse));

return IcebergRestUtils.ok(loadTableResponse);
}

@POST
Expand All @@ -122,10 +146,17 @@ public Response updateTable(
}
TableIdentifier tableIdentifier =
TableIdentifier.of(RESTUtil.decodeNamespace(namespace), table);
return IcebergRestUtils.ok(
eventBus.dispatchPreEvent(
new IcebergUpdateTablePreEvent(
"user", NameIdentifier.of(namespace, table), updateTableRequest));
LoadTableResponse loadTableResponse =
icebergCatalogWrapperManager
.getOps(prefix)
.updateTable(tableIdentifier, updateTableRequest));
.updateTable(tableIdentifier, updateTableRequest);
eventBus.dispatchEvent(
new IcebergUpdateTablePostEvent(
"user", NameIdentifier.of(namespace, table), updateTableRequest, loadTableResponse));
return IcebergRestUtils.ok(loadTableResponse);
}

@DELETE
Expand Down
Loading
Loading