Advertisement
Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- package de.hpi.ddm.actors;
- // custom imports
- import akka.Done;
- import akka.stream.javadsl.*;
- import com.google.common.primitives.Bytes;
- import java.util.List;
- import java.io.Serializable;
- import java.util.Arrays;
- import java.util.concurrent.CompletionStage;
- import akka.NotUsed;
- import akka.actor.AbstractLoggingActor;
- import akka.actor.ActorRef;
- import akka.actor.ActorSelection;
- import akka.actor.Props;
- import akka.serialization.Serialization;
- import akka.serialization.SerializationExtension;
- import akka.serialization.Serializers;
- import lombok.AllArgsConstructor;
- import lombok.Data;
- import lombok.NoArgsConstructor;
- public class LargeMessageProxy extends AbstractLoggingActor {
- ////////////////////////
- // Actor Construction //
- ////////////////////////
- public static final String DEFAULT_NAME = "largeMessageProxy";
- public static Props props() {
- return Props.create(LargeMessageProxy.class);
- }
- ////////////////////
- // Actor Messages //
- ////////////////////
- @Data @NoArgsConstructor @AllArgsConstructor
- public static class LargeMessage<T> implements Serializable {
- private static final long serialVersionUID = 2940665245810221108L;
- private T message;
- private ActorRef receiver;
- }
- @Data @NoArgsConstructor @AllArgsConstructor
- public static class BytesMessage<T> implements Serializable {
- private static final long serialVersionUID = 4057807743872319842L;
- private T bytes;
- private ActorRef sender;
- private ActorRef receiver;
- private Integer serializerID;
- private String manifest;
- private RunnableGraph runnable;
- public T getBytes() {
- return bytes;
- }
- public ActorRef getReceiver() {
- return receiver;
- }
- public ActorRef getSender() {
- return sender;
- }
- public Integer getSerializerID() {
- return serializerID;
- }
- public String getManifest() {
- return manifest;
- }
- public RunnableGraph getRunnable(){
- return runnable;
- }
- }
- /////////////////
- // Actor State //
- /////////////////
- /////////////////////
- // Actor Lifecycle //
- /////////////////////
- ////////////////////
- // Actor Behavior //
- ////////////////////
- @Override
- public Receive createReceive() {
- return receiveBuilder()
- .match(LargeMessage.class, this::handle)
- .match(BytesMessage.class, this::handle)
- .matchAny(object -> this.log().info("Received unknown message: \"{}\"", object.toString()))
- .build();
- }
- private void handle(LargeMessage<?> message) {
- ActorRef receiver = message.getReceiver();
- ActorSelection receiverProxy = this.context().actorSelection(receiver.path().child(DEFAULT_NAME));
- // This will definitely fail in a distributed setting if the serialized message is large!
- // Solution options:
- // 1. Serialize the object and send its bytes batch-wise (make sure to use artery's side channel then).
- // 2. Serialize the object and send its bytes via Akka streaming.
- // 3. Send the object via Akka's http client-server component.
- // 4. Other ideas ...
- // Here we serialize our data
- Serialization serialization = SerializationExtension.get(getContext().getSystem());
- byte[] bytes = serialization.serialize(message).get();
- String bytesAsString = new String(bytes);
- // List<Byte> iterableByteArray = Bytes.asList(bytes); // this one can be an option to make a byte array iterable
- int serializerId = serialization.findSerializerFor(message).identifier();
- String manifest = Serializers.manifestFor(serialization.findSerializerFor(message), message);
- final Source<String, NotUsed> source = Source.from(Arrays.asList(bytesAsString));
- final Flow<String, String, NotUsed> flow = Flow.fromFunction((String x) ->x);
- final Sink<String, CompletionStage<Done>> sink = Sink.foreach(x -> System.out.println(x));
- final RunnableGraph<NotUsed> runnable = source.via(flow).to(sink);
- receiverProxy.tell(new BytesMessage<>(message.getMessage(), this.sender(), message.getReceiver(), serializerId, manifest, runnable), this.self());
- }
- private void handle(BytesMessage<?> message) {
- // Reassemble the message content, deserialize it and/or load the content from some local location before forwarding its content.
- Integer serializerId = message.serializerID;
- String manifest = message.manifest;
- RunnableGraph runnable = message.runnable;
- runnable.run();
- message.getReceiver().tell(message.getBytes(), message.getSender());
- }
- }
Advertisement
Add Comment
Please, Sign In to add comment
Advertisement