-
Notifications
You must be signed in to change notification settings - Fork 495
AWS CloudWatch Event Sink Implementation #1965
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
Changes from 43 commits
2e07dde
06eaca4
1515fbb
854501b
c1c94b2
ab3c5f9
ab9ccbe
a641136
5a355d1
bab4439
04c310a
d4b44ff
4d0554a
cc715ad
518aaaa
8758255
f3f62a0
d21dabc
9054511
828760a
ae79600
9d47684
025de74
e4ec3f8
491ea3a
d453660
f89b0ae
e5c02b7
4f8a15b
1305321
27f28f4
6b42071
ec2bee8
69b7feb
e8b5e93
3030d6a
b9abab6
4c91c57
b0c6160
44fad9a
e8447a9
688ec97
360cb91
7c09d9b
5e34884
e2ed743
a870aea
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,44 @@ | ||
| /* | ||
| * 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.polaris.service.events.jsonEventListener; | ||
|
|
||
| import java.util.HashMap; | ||
| import org.apache.polaris.service.events.AfterTableRefreshedEvent; | ||
| import org.apache.polaris.service.events.PolarisEventListener; | ||
|
|
||
| /** | ||
| * Abstract base class from which all event sinks that output events in JSON format can extend. | ||
| * | ||
| * <p>This class provides a common framework for transforming Polaris events into JSON format and | ||
| * sending them to various destinations. Concrete implementations should override the {@link | ||
| * #transformAndSendEvent(HashMap)} method to define how the JSON event data should be transmitted | ||
| * or stored. | ||
| */ | ||
| public abstract class JsonEventListener extends PolarisEventListener { | ||
|
eric-maynard marked this conversation as resolved.
Outdated
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. A note, that I really like having this as a public class because if someone does want to write a plugin that doesn't interact with any of the Iceberg or other based libraries, they can always utilize this instead and only depend on Polaris. |
||
| protected abstract void transformAndSendEvent(HashMap<String, Object> properties); | ||
|
|
||
| @Override | ||
| public void onAfterTableRefreshed(AfterTableRefreshedEvent event) { | ||
| HashMap<String, Object> properties = new HashMap<>(); | ||
| properties.put("event_type", event.getClass().getSimpleName()); | ||
| properties.put("table_identifier", event.tableIdentifier().toString()); | ||
| transformAndSendEvent(properties); | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,31 @@ | ||
| /* | ||
| * 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.polaris.service.events.jsonEventListener.aws.cloudwatch; | ||
|
|
||
| /** Configuration interface for AWS CloudWatch event listener settings. */ | ||
| public interface AwsCloudWatchConfiguration { | ||
| String awsCloudwatchlogGroup(); | ||
|
|
||
| String awsCloudwatchlogStream(); | ||
|
|
||
| String awsCloudwatchRegion(); | ||
|
|
||
| boolean synchronousMode(); | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,158 @@ | ||
| /* | ||
| * 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.polaris.service.events.jsonEventListener.aws.cloudwatch; | ||
|
|
||
| import com.fasterxml.jackson.core.JsonProcessingException; | ||
| import com.fasterxml.jackson.databind.ObjectMapper; | ||
| import io.smallrye.common.annotation.Identifier; | ||
| import jakarta.annotation.PostConstruct; | ||
| import jakarta.annotation.PreDestroy; | ||
| import jakarta.enterprise.context.ApplicationScoped; | ||
| import jakarta.inject.Inject; | ||
| import jakarta.ws.rs.core.Context; | ||
| import jakarta.ws.rs.core.SecurityContext; | ||
| import java.time.Clock; | ||
| import java.util.HashMap; | ||
| import java.util.List; | ||
| import java.util.concurrent.CompletableFuture; | ||
| import java.util.concurrent.CompletionException; | ||
| import org.apache.polaris.core.context.CallContext; | ||
| import org.apache.polaris.service.events.jsonEventListener.JsonEventListener; | ||
| import org.slf4j.Logger; | ||
| import org.slf4j.LoggerFactory; | ||
| import software.amazon.awssdk.regions.Region; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.CloudWatchLogsAsyncClient; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.CreateLogGroupRequest; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.CreateLogGroupResponse; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.CreateLogStreamRequest; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.CreateLogStreamResponse; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.InputLogEvent; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.PutLogEventsRequest; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.PutLogEventsResponse; | ||
| import software.amazon.awssdk.services.cloudwatchlogs.model.ResourceAlreadyExistsException; | ||
|
|
||
| @ApplicationScoped | ||
| @Identifier("aws-cloudwatch") | ||
| public class AwsCloudWatchEventListener extends JsonEventListener { | ||
| private static final Logger LOGGER = LoggerFactory.getLogger(AwsCloudWatchEventListener.class); | ||
| private final ObjectMapper objectMapper = new ObjectMapper(); | ||
|
singhpk234 marked this conversation as resolved.
Outdated
|
||
|
|
||
| private CloudWatchLogsAsyncClient client; | ||
|
|
||
| private final String logGroup; | ||
| private final String logStream; | ||
| private final Region region; | ||
| private final boolean synchronousMode; | ||
| private final Clock clock; | ||
|
|
||
| @Inject CallContext callContext; | ||
|
|
||
| @Context SecurityContext securityContext; | ||
|
|
||
| @Inject | ||
| public AwsCloudWatchEventListener(AwsCloudWatchConfiguration config, Clock clock) { | ||
| this.logStream = config.awsCloudwatchlogStream(); | ||
| this.logGroup = config.awsCloudwatchlogGroup(); | ||
| this.region = Region.of(config.awsCloudwatchRegion()); | ||
| this.synchronousMode = config.synchronousMode(); | ||
| this.clock = clock; | ||
| } | ||
|
|
||
| @PostConstruct | ||
| void start() { | ||
| this.client = createCloudWatchAsyncClient(); | ||
| ensureLogGroupAndStream(); | ||
| } | ||
|
|
||
| protected CloudWatchLogsAsyncClient createCloudWatchAsyncClient() { | ||
| return CloudWatchLogsAsyncClient.builder().region(region).build(); | ||
| } | ||
|
|
||
| private void ensureLogGroupAndStream() { | ||
| try { | ||
| CompletableFuture<CreateLogGroupResponse> future = | ||
| client.createLogGroup(CreateLogGroupRequest.builder().logGroupName(logGroup).build()); | ||
| future.join(); | ||
|
adnanhemani marked this conversation as resolved.
Outdated
|
||
| } catch (CompletionException e) { | ||
| if (e.getCause() instanceof ResourceAlreadyExistsException) { | ||
| LOGGER.debug("Log group {} already exists", logGroup); | ||
| } else { | ||
| throw e; | ||
| } | ||
| } | ||
|
|
||
| try { | ||
| CompletableFuture<CreateLogStreamResponse> future = | ||
| client.createLogStream( | ||
| CreateLogStreamRequest.builder() | ||
| .logGroupName(logGroup) | ||
| .logStreamName(logStream) | ||
| .build()); | ||
| future.join(); | ||
| } catch (CompletionException e) { | ||
| if (e.getCause() instanceof ResourceAlreadyExistsException) { | ||
| LOGGER.debug("Log stream {} already exists", logStream); | ||
| } else { | ||
| throw e; | ||
| } | ||
| } | ||
| } | ||
|
|
||
| @PreDestroy | ||
| void shutdown() { | ||
| if (client != null) { | ||
| client.close(); | ||
|
singhpk234 marked this conversation as resolved.
|
||
| } | ||
| } | ||
|
|
||
| @Override | ||
| protected void transformAndSendEvent(HashMap<String, Object> properties) { | ||
| properties.put("realm", callContext.getRealmContext().getRealmIdentifier()); | ||
|
singhpk234 marked this conversation as resolved.
Outdated
|
||
| properties.put("principal", securityContext.getUserPrincipal().getName()); | ||
|
singhpk234 marked this conversation as resolved.
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit but I would recommend putting these into constants somewhere, e.g.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Putting them as class-wide constants isn't a great idea IMO. These variables are RequestScoped, while the class as a whole is ApplicationScoped. The best we can probably do is to keep them as variables within the
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. To be clear, I'm talking about the string literals like |
||
| // TODO: Add request ID when it is available | ||
| String eventAsJson; | ||
| try { | ||
| eventAsJson = objectMapper.writeValueAsString(properties); | ||
| } catch (JsonProcessingException e) { | ||
| LOGGER.error("Error processing event into JSON string: ", e); | ||
|
singhpk234 marked this conversation as resolved.
|
||
| return; | ||
| } | ||
| InputLogEvent inputLogEvent = | ||
| InputLogEvent.builder().message(eventAsJson).timestamp(clock.millis()).build(); | ||
| PutLogEventsRequest.Builder requestBuilder = | ||
| PutLogEventsRequest.builder() | ||
| .logGroupName(logGroup) | ||
| .logStreamName(logStream) | ||
| .logEvents(List.of(inputLogEvent)); | ||
| CompletableFuture<PutLogEventsResponse> future = | ||
| client | ||
| .putLogEvents(requestBuilder.build()) | ||
| .whenComplete( | ||
| (resp, err) -> { | ||
| if (err != null) { | ||
| LOGGER.error( | ||
| "Error writing log to CloudWatch. Event: {}, Error: ", inputLogEvent, err); | ||
| } | ||
| }); | ||
| if (synchronousMode) { | ||
| future.join(); | ||
| } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,99 @@ | ||
| /* | ||
| * 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.polaris.service.quarkus.events.jsonEventListener.aws.cloudwatch; | ||
|
|
||
| import io.quarkus.runtime.annotations.StaticInitSafe; | ||
| import io.smallrye.config.ConfigMapping; | ||
| import io.smallrye.config.WithDefault; | ||
| import io.smallrye.config.WithName; | ||
| import jakarta.enterprise.context.ApplicationScoped; | ||
| import org.apache.polaris.service.events.jsonEventListener.aws.cloudwatch.AwsCloudWatchConfiguration; | ||
|
|
||
| /** | ||
| * Quarkus-specific configuration interface for AWS CloudWatch event listener integration. | ||
| * | ||
| * <p>This interface extends the base {@link AwsCloudWatchConfiguration} and provides | ||
| * Quarkus-specific configuration mappings for AWS CloudWatch logging functionality. | ||
| */ | ||
| @StaticInitSafe | ||
| @ConfigMapping(prefix = "polaris.event-listener.aws-cloudwatch") | ||
| @ApplicationScoped | ||
| public interface QuarkusAwsCloudWatchConfiguration extends AwsCloudWatchConfiguration { | ||
|
eric-maynard marked this conversation as resolved.
|
||
|
|
||
| /** | ||
| * Returns the AWS CloudWatch log group name for event logging. | ||
| * | ||
| * <p>The log group is a collection of log streams that share the same retention, monitoring, and | ||
| * access control settings. If not specified, defaults to "polaris-cloudwatch-default-group". | ||
| * | ||
| * <p>Configuration property: {@code polaris.event-listener.aws-cloudwatch.log-group} | ||
| * | ||
| * @return a String containing the log group name, or the default value if not configured | ||
| */ | ||
| @WithName("log-group") | ||
| @WithDefault("polaris-cloudwatch-default-group") | ||
| @Override | ||
| String awsCloudwatchlogGroup(); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Good point, changed. |
||
|
|
||
| /** | ||
| * Returns the AWS CloudWatch log stream name for event logging. | ||
| * | ||
| * <p>A log stream is a sequence of log events that share the same source. Each log stream belongs | ||
| * to one log group. If not specified, defaults to "polaris-cloudwatch-default-stream". | ||
| * | ||
| * <p>Configuration property: {@code polaris.event-listener.aws-cloudwatch.log-stream} | ||
| * | ||
| * @return a String containing the log stream name, or the default value if not configured | ||
| */ | ||
| @WithName("log-stream") | ||
| @WithDefault("polaris-cloudwatch-default-stream") | ||
| @Override | ||
| String awsCloudwatchlogStream(); | ||
|
|
||
| /** | ||
| * Returns the AWS region where CloudWatch logs should be sent. | ||
| * | ||
| * <p>This specifies the AWS region for the CloudWatch service endpoint. The region must be a | ||
| * valid AWS region identifier. If not specified, defaults to "us-east-1". | ||
| * | ||
| * <p>Configuration property: {@code polaris.event-listener.aws-cloudwatch.region} | ||
| * | ||
| * @return a String containing the AWS region, or the default value if not configured | ||
| */ | ||
| @WithName("region") | ||
| @WithDefault("us-east-1") | ||
| @Override | ||
| String awsCloudwatchRegion(); | ||
|
adnanhemani marked this conversation as resolved.
Outdated
|
||
|
|
||
| /** | ||
| * Returns the synchronous mode setting for CloudWatch logging. | ||
| * | ||
| * <p>When set to "true", log events are sent to CloudWatch synchronously, which may impact | ||
| * application performance but ensures immediate delivery. When set to "false" (default), log | ||
| * events are sent asynchronously for better performance. | ||
| * | ||
| * <p>Configuration property: {@code polaris.event-listener.aws-cloudwatch.synchronous-mode} | ||
| * | ||
| * @return a boolean value indicating the synchronous mode setting | ||
| */ | ||
| @WithName("synchronous-mode") | ||
| @WithDefault("false") | ||
| @Override | ||
| boolean synchronousMode(); | ||
| } | ||
Uh oh!
There was an error while loading. Please reload this page.