Click here to close now.

Welcome!

Java Authors: Liz McMillan, Carmen Gonzalez, Pat Romanski, Jason Bloomberg, Roger Strukhoff

Related Topics: Java

Java: Article

Next-Gen Concurrency in Java: The Actor Model

Multiple concurrent processes can communicate with each other without needing to use shared state variables

At a time where the clock speeds of processors have been stable over the past couple of years, and Moore's Law is instead being applied by increasing the number of processor cores, it is getting more important for applications to use concurrent processing to reduce run/response times, as the time slicing routine via increased clock speed will no longer be available to bail out slow running programs.

Carl Hewitt proposed the Actor Model in 1973 as a way to implement unbounded nondeterminism in concurrent processing. In many ways this model was influenced by the packet switching mechanism, for example, no synchronous handshake between sender and receiver, inherently concurrent message passing, messages may not arrive in the order they were sent, addresses are stored in messages, etc.

The main difference between this model and most other parallel processing systems is that it uses message passing instead of shared variables for communication between concurrent processes. Using shared memory to communicate between concurrent processes requires the application of some form of locking mechanism to coordinate between threads, which may give rise to live locks, deadlocks, race conditions and starvation.

Actors are the location transparent primitives that form the basis of the actor model. Each actor is associated with a mailbox (which is a queue with multiple producers and a single consumer) where it receives and buffers messages, and a behavior that is executed as a reaction to a message received. The messages are immutable and may be passed between actors synchronously or asynchronously depending on the type of operation being invoked. In response to a message that it receives, an actor can make local decisions, create more actors, send more messages, and designate how to respond to the next message received. Actors never share state and thus don't need to compete for locks for access to shared data.

The actor model first rose to fame in the language Erlang, designed by Ericcson in 1986. It has since been implemented in many next-generation languages on the JVM such as Scala, Groovy and Fantom. It is the simplicity of usage provided via a higher level of abstraction that makes the actor model easier to implement and reason about.

It's now possible to implement the actor model in Java, thanks to the growing number of third-party concurrency libraries advertising this feature. Akka is one such library, written in Scala, that uses the Actor model to simplify writing fault-tolerant, highly scalable applications in both Java and Scala.

Implementation
Using Akka, we shall attempt to create a concurrent processing system for loan request processing in a bank as can be seen in the Figure 1.

Figure 1

The system consists of four actors:

  1. The front desk - which shall receive loan requests from the customers and send them to the back office for processing. It shall also maintain the statistics of the number of loans accepted/rejected and print a report detailing the same on being asked to do so.
  2. The back office - which shall sort the loan requests into personal loans and home loans, and send them to the corresponding accountant for approval/rejection.
  3. The personal loan accountant - who shall process personal loan requests, including approving/rejecting the requests, carrying out credit history checks and calculating the rate of interest.
  4. The home loan accountant - who shall process home loan requests, including approving/rejecting the requests, carrying out credit history checks and calculating the rate of interest.

To get started, download the version 2.0 of Akka for Java from http://akka.io/downloads and add the jars present in ‘akka-2.0-RC4\lib' to the classpath.

Next, create a new Java class "Bank.java" and add the following import statements :

import akka.actor.ActorRef;

import akka.actor.ActorSystem;

import akka.actor.Props;

import akka.actor.UntypedActor;

import akka.routing.RoundRobinRouter;

Now, create a few static nested classes under ‘Bank' that will act as messages (DTOs) and be passed to actors.

1. ‘LoanRequest' will contain the following elements and their corresponding getters/setters:

int requestedLoan;
int
accountBalance;

2. ‘PersonalLoanRequest' will extend ‘LoanRequest' and contain the following element and its corresponding getter:

final static String type="Personal";

3. ‘HomeLoanRequest' will extend ‘LoanRequest' and contain the following element and its corresponding getter:

final static String type="Home";

4. ‘LoanReply' will contain the following elements and their corresponding getters/setters:

String type;
boolean
approved;
int
rate;

Next, create a static nested class ‘PersonalLoanAccountant' under ‘Bank'.

public static class PersonalLoanAccountant extends UntypedActor {
public
int rateCaluclation(int requestedLoan, int accountBalance) {
if
(accountBalance/requestedLoan>=2)
return
5;
else

return
6;
}

public
void checkCreditHistory() {
for
(int i=0; i<1000; i++) {
continue
;
}
}
public
void onReceive(Object message) {
if
(message instanceof PersonalLoanRequest) {
PersonalLoanRequest request=PersonalLoanRequest.class.cast(message);
LoanReply reply = new LoanReply();
reply.setType(request.getType());
if
(request.getRequestedLoan()<request.getaccountBalance()) {
reply.setApproved(true);
reply.setRate(rateCaluclation(request.getRequestedLoan(), request.getaccountBalance()));
checkCreditHistory();
}
getSender().tell(reply);
} else {
unhandled(message);
}
}
}

The above class  serves as an actor in the system. It extends the ‘UntypedActor' base class provided by Akka and must define its ‘onReceive‘ method. This method acts as the mailbox and receives messages from other actors (or non-actors) in the system. If the message received is of type 'PersonalLoanRequest', then it can be processed by approving/rejecting the loan request, setting the rate of interest and checking the requestor's credit history. Once the processing is complete, the requestor uses the sender's reference (which is embedded in the message) to send the reply (LoanReply) to the requestor via the ‘getSender().tell()' method.

Now, create a similar static nested class ‘HomeLoanAccountant' under ‘Bank'

public static class HomeLoanAccountant extends UntypedActor {
public
int rateCaluclation(int requestedLoan, int accountBalance) {
if
(accountBalance/requestedLoan>=2)
return
7;
else

return
8;
}

public
void checkCreditHistory() {
for
(int i=0; i<2000; i++) {
continue
;
}
}
public
void onReceive(Object message) {
if
(message instanceof HomeLoanRequest) {
HomeLoanRequest request=HomeLoanRequest.class.cast(message);
LoanReply reply = new LoanReply();
reply.setType(request.getType());
if
(request.getRequestedLoan()<request.getaccountBalance()) {
reply.setApproved(true);
reply.setRate(rateCaluclation(request.getRequestedLoan(), request.getaccountBalance()));
checkCreditHistory();
}
getSender().tell(reply);
} else {
unhandled(message);
}
}
}

Now, create a static nested class ‘BackOffice' under ‘Bank':

public static class BackOffice extends UntypedActor {

ActorRef personalLoanAccountant=getContext().actorOf(new Props(PersonalLoanAccountant.class).withRouter(new RoundRobinRouter(2)));
ActorRef homeLoanAccountant=getContext().actorOf(new Props(HomeLoanAccountant.class).withRouter(new RoundRobinRouter(2)));

public void onReceive(Object message) {
if
(message instanceof PersonalLoanRequest) {
personalLoanAccountant.forward(message, getContext());
} else if (message instanceof HomeLoanRequest) {
homeLoanAccountant.forward(message, getContext());
} else {
unhandled(message);
}
}

}

The above class serves as an actor in the system and is purely used to route the incoming messages to PersonalLoanAccountant or HomeLoanAccountant, based on the type of message received. It defines actor references ‘personalLoanAccountant' and ‘homeLoanAccountant' to ‘PersonalLoanAccountant.class' and ‘HomeLoanAccountant.class', respectively. Each of these references initiates two instances of the actor it refers to and attaches a round-robin router to cycle through the actor instances. The ‘onReceive' method checks the type of the message received and forwards the message to either ‘PersonalLoanAccountant' or ‘HomeLoanAccountant' based on the message type. The ‘forward()' method helps ensure that the reference of the original sender is maintained in the message, so the receiver of the message (‘PersonalLoanAccountant' or ‘HomeLoanAccountant') can directly reply back to the message's original sender.

Now, create a static nested class ‘FrontDesk' under ‘Bank':

public static class FrontDesk extends UntypedActor {
int
approvedPersonalLoans=0;
int
approvedHomeLoans=0;
int
rejectedPersonalLoans=0;
int
rejectedHomeLoans=0;

ActorRef backOffice=getContext().actorOf(
new Props(BackOffice.class), "backOffice");

public
void maintainLoanApprovalStats(Object message) {
LoanReply reply = LoanReply.class.cast(message);
if
(reply.isApproved()) {
System.out.println(reply.getType()+" Loan Approved"+" at "+reply.getRate()+"% interest.");
if
(reply.getType().equals("Personal"))
++approvedPersonalLoans;
else
if(reply.getType().equals("Home"))
++approvedHomeLoans;
} else {
System.out.println(reply.getType()+" Loan Rejected");
if
(reply.getType().equals("Personal"))
++rejectedPersonalLoans;
else
if(reply.getType().equals("Home"))
++rejectedHomeLoans;
}
}

public void printLoanApprovalStats() {
System.out.println("--- REPORT ---");
System.out.println("Personal Loans Approved : "+approvedPersonalLoans);
System.out.println("Home Loans Approved : "+approvedHomeLoans);
System.out.println("Personal Loans Rejected : "+rejectedPersonalLoans);
System.out.println("Home Loans Rejected : "+rejectedHomeLoans);
}
public
void onReceive(Object message) {
if
(message instanceof LoanRequest) {
backOffice.tell(message, getSelf());
} else if (message instanceof LoanReply) {
maintainLoanApprovalStats(message);
} else if(message instanceof String && message.equals("printLoanApprovalStats")) {
printLoanApprovalStats();
getContext().stop(getSelf());
} else {
unhandled(message);
}
}
}

The above class serves as the final actor in the system. It creates a reference ‘backOffice' to the actor ‘BackOffice.class'. If the message received is of type ‘LoanRequest', it sends the message to ‘BackOffice' via the method ‘backOffice.tell()', which takes a message and the reference to the sender (acquired through the method ‘getSelf()') as arguments. If the message received is of type ‘LoanReply', it updates the counters to maintain the approved/rejected counts. If the message received is "printLoanApprovalStats", it prints the stats stored in the counters, and then proceeds to stop itself via the method ‘getContext().stop(getSelf())'. The actors follow a pattern of supervisor hierarchy, and thus this command trickles down the hierarchy chain and stops all four actors in the system.

Finally, write a few methods under ‘Bank' to submit requests to the ‘FrontOffice':

public static void main(String[] args) throws InterruptedException {
ActorSystem system = ActorSystem.create("bankSystem");
ActorRef frontDesk = system.actorOf(new Props(FrontDesk.class), "frontDesk");
submitLoanRequests
(frontDesk);
Thread.sleep(1000);
printLoanApprovalStats
(frontDesk);
system.shutdown();
}

public
static PersonalLoanRequest getPersonalLoanRequest() {
int
min=10000; int max=50000;
int
amount=min + (int)(Math.random() * ((max - min) + 1));
int
balance=min + (int)(Math.random() * ((max - min) + 1));
return
(new Bank()).new PersonalLoanRequest(amount, balance);
}

public
static HomeLoanRequest getHomeLoanRequest() {
int
min=50000; int max=90000;
int
amount=min + (int)(Math.random() * ((max - min) + 1));
int
balance=min + (int)(Math.random() * ((max - min) + 1));
return
(new Bank()).new HomeLoanRequest(amount, balance);
}

public static void submitLoanRequests(ActorRef frontDesk) {
for
(int i=0;i<1000;i++) {
frontDesk.tell(getPersonalLoanRequest());
frontDesk.tell(getHomeLoanRequest());
}
}

public static void printLoanApprovalStats(ActorRef frontDesk) {
frontDesk.tell("printLoanApprovalStats");
}

The method ‘main' creates an actor system ‘system' using the method ‘ActorSystem.create()'. It then creates a reference ‘frontDesk' to the actor ‘FrontDesk.class'. It uses this reference to send a 1000 requests each of types ‘PersonalLoanRequest' and ‘HomeLoanRequest' to ‘FrontDesk'. It then sleeps for a second, following which it sends the message "printLoanApprovalStats" to ‘FrontDesk'. Once done, it shuts down the actor system via the method ‘system.shutdown()'.

Conclusion
Running the above code will create a loan request processing system with concurrent processing capabilities, and all this without using a single synchronize/lock pattern. Moreover, the code doesn't need one to go into the low-level semantics of the JVM threading mechanism or use the complex ‘java.util.concurrent' package. This mechanism of concurrent processing using the Actor Model is truly more robust as multiple concurrent processes can communicate with each other without needing to use shared state variables.

More Stories By Sanat Vij

Sanat Vij is a professional software engineer currently working at CenturyLink. He has vast experience in developing high availability applications, configuring application servers, JVM profiling and memory management. He specializes in performance tuning of applications, reducing response times, and increasing stability.

Comments (0)

Share your thoughts on this story.

Add your comment
You must be signed in to add a comment. Sign-in | Register

In accordance with our Comment Policy, we encourage comments that are on topic, relevant and to-the-point. We will remove comments that include profanity, personal attacks, racial slurs, threats of violence, or other inappropriate material that violates our Terms and Conditions, and will block users who make repeated violations. We ask all readers to expect diversity of opinion and to treat one another with dignity and respect.


@ThingsExpo Stories
The Workspace-as-a-Service (WaaS) market will grow to $6.4B by 2018. In his session at 16th Cloud Expo, Seth Bostock, CEO of IndependenceIT, will begin by walking the audience through the evolution of Workspace as-a-Service, where it is now vs. where it going. To look beyond the desktop we must understand exactly what WaaS is, who the users are, and where it is going in the future. IT departments, ISVs and service providers must look to workflow and automation capabilities to adapt to growing demand and the rapidly changing workspace model.
Even as cloud and managed services grow increasingly central to business strategy and performance, challenges remain. The biggest sticking point for companies seeking to capitalize on the cloud is data security. Keeping data safe is an issue in any computing environment, and it has been a focus since the earliest days of the cloud revolution. Understandably so: a lot can go wrong when you allow valuable information to live outside the firewall. Recent revelations about government snooping, along with a steady stream of well-publicized data breaches, only add to the uncertainty
SYS-CON Events announced today that Dyn, the worldwide leader in Internet Performance, will exhibit at SYS-CON's 16th International Cloud Expo®, which will take place on June 9-11, 2015, at the Javits Center in New York City, NY. Dyn is a cloud-based Internet Performance company. Dyn helps companies monitor, control, and optimize online infrastructure for an exceptional end-user experience. Through a world-class network and unrivaled, objective intelligence into Internet conditions, Dyn ensures traffic gets delivered faster, safer, and more reliably than ever.
Hadoop as a Service (as offered by handful of niche vendors now) is a cloud computing solution that makes medium and large-scale data processing accessible, easy, fast and inexpensive. In his session at Big Data Expo, Kumar Ramamurthy, Vice President and Chief Technologist, EIM & Big Data, at Virtusa, will discuss how this is achieved by eliminating the operational challenges of running Hadoop, so one can focus on business growth. The fragmented Hadoop distribution world and various PaaS solutions that provide a Hadoop flavor either make choices for customers very flexible in the name of opti...
As organizations shift toward IT-as-a-service models, the need for managing and protecting data residing across physical, virtual, and now cloud environments grows with it. CommVault can ensure protection &E-Discovery of your data – whether in a private cloud, a Service Provider delivered public cloud, or a hybrid cloud environment – across the heterogeneous enterprise. In his session at 16th Cloud Expo, Randy De Meno, Chief Technologist - Windows Products and Microsoft Partnerships, will discuss how to cut costs, scale easily, and unleash insight with CommVault Simpana software, the only si...
Cloud data governance was previously an avoided function when cloud deployments were relatively small. With the rapid adoption in public cloud – both rogue and sanctioned, it’s not uncommon to find regulated data dumped into public cloud and unprotected. This is why enterprises and cloud providers alike need to embrace a cloud data governance function and map policies, processes and technology controls accordingly. In her session at 15th Cloud Expo, Evelyn de Souza, Data Privacy and Compliance Strategy Leader at Cisco Systems, will focus on how to set up a cloud data governance program and s...
Roberto Medrano, Executive Vice President at SOA Software, had reached 30,000 page views on his home page - http://RobertoMedrano.SYS-CON.com/ - on the SYS-CON family of online magazines, which includes Cloud Computing Journal, Internet of Things Journal, Big Data Journal, and SOA World Magazine. He is a recognized executive in the information technology fields of SOA, internet security, governance, and compliance. He has extensive experience with both start-ups and large companies, having been involved at the beginning of four IT industries: EDA, Open Systems, Computer Security and now SOA.
The industrial software market has treated data with the mentality of “collect everything now, worry about how to use it later.” We now find ourselves buried in data, with the pervasive connectivity of the (Industrial) Internet of Things only piling on more numbers. There’s too much data and not enough information. In his session at @ThingsExpo, Bob Gates, Global Marketing Director, GE’s Intelligent Platforms business, to discuss how realizing the power of IoT, software developers are now focused on understanding how industrial data can create intelligence for industrial operations. Imagine ...
Operational Hadoop and the Lambda Architecture for Streaming Data Apache Hadoop is emerging as a distributed platform for handling large and fast incoming streams of data. Predictive maintenance, supply chain optimization, and Internet-of-Things analysis are examples where Hadoop provides the scalable storage, processing, and analytics platform to gain meaningful insights from granular data that is typically only valuable from a large-scale, aggregate view. One architecture useful for capturing and analyzing streaming data is the Lambda Architecture, representing a model of how to analyze rea...
SYS-CON Events announced today that Vitria Technology, Inc. will exhibit at SYS-CON’s @ThingsExpo, which will take place on June 9-11, 2015, at the Javits Center in New York City, NY. Vitria will showcase the company’s new IoT Analytics Platform through live demonstrations at booth #330. Vitria’s IoT Analytics Platform, fully integrated and powered by an operational intelligence engine, enables customers to rapidly build and operationalize advanced analytics to deliver timely business outcomes for use cases across the industrial, enterprise, and consumer segments.
The Internet of Things (IoT) promises to evolve the way the world does business; however, understanding how to apply it to your company can be a mystery. Most people struggle with understanding the potential business uses or tend to get caught up in the technology, resulting in solutions that fail to meet even minimum business goals. In his session at @ThingsExpo, Jesse Shiah, CEO / President / Co-Founder of AgilePoint Inc., showed what is needed to leverage the IoT to transform your business. He discussed opportunities and challenges ahead for the IoT from a market and technical point of vie...
Advanced Persistent Threats (APTs) are increasing at an unprecedented rate. The threat landscape of today is drastically different than just a few years ago. Attacks are much more organized and sophisticated. They are harder to detect and even harder to anticipate. In the foreseeable future it's going to get a whole lot harder. Everything you know today will change. Keeping up with this changing landscape is already a daunting task. Your organization needs to use the latest tools, methods and expertise to guard against those threats. But will that be enough? In the foreseeable future attacks w...
HP and Aruba Networks on Monday announced a definitive agreement for HP to acquire Aruba, a provider of next-generation network access solutions for the mobile enterprise, for $24.67 per share in cash. The equity value of the transaction is approximately $3.0 billion, and net of cash and debt approximately $2.7 billion. Both companies' boards of directors have approved the deal. "Enterprises are facing a mobile-first world and are looking for solutions that help them transition legacy investments to the new style of IT," said Meg Whitman, Chairman, President and Chief Executive Officer of HP...
Containers and microservices have become topics of intense interest throughout the cloud developer and enterprise IT communities. Accordingly, attendees at the upcoming 16th Cloud Expo at the Javits Center in New York June 9-11 will find fresh new content in a new track called PaaS | Containers & Microservices Containers are not being considered for the first time by the cloud community, but a current era of re-consideration has pushed them to the top of the cloud agenda. With the launch of Docker's initial release in March of 2013, interest was revved up several notches. Then late last...
Disruptive macro trends in technology are impacting and dramatically changing the "art of the possible" relative to supply chain management practices through the innovative use of IoT, cloud, machine learning and Big Data to enable connected ecosystems of engagement. Enterprise informatics can now move beyond point solutions that merely monitor the past and implement integrated enterprise fabrics that enable end-to-end supply chain visibility to improve customer service delivery and optimize supplier management. Learn about enterprise architecture strategies for designing connected systems tha...
The explosion of connected devices / sensors is creating an ever-expanding set of new and valuable data. In parallel the emerging capability of Big Data technologies to store, access, analyze, and react to this data is producing changes in business models under the umbrella of the Internet of Things (IoT). In particular within the Insurance industry, IoT appears positioned to enable deep changes by altering relationships between insurers, distributors, and the insured. In his session at @ThingsExpo, Michael Sick, a Senior Manager and Big Data Architect within Ernst and Young's Financial Servi...
The explosion of connected devices / sensors is creating an ever-expanding set of new and valuable data. In parallel the emerging capability of Big Data technologies to store, access, analyze, and react to this data is producing changes in business models under the umbrella of the Internet of Things (IoT). In particular within the Insurance industry, IoT appears positioned to enable deep changes by altering relationships between insurers, distributors, and the insured. In his session at @ThingsExpo, Michael Sick, a Senior Manager and Big Data Architect within Ernst and Young's Financial Servi...
PubNub on Monday has announced that it is partnering with IBM to bring its sophisticated real-time data streaming and messaging capabilities to Bluemix, IBM’s cloud development platform. “Today’s app and connected devices require an always-on connection, but building a secure, scalable solution from the ground up is time consuming, resource intensive, and error-prone,” said Todd Greene, CEO of PubNub. “PubNub enables web, mobile and IoT developers building apps on IBM Bluemix to quickly add scalable realtime functionality with minimal effort and cost.”
Sensor-enabled things are becoming more commonplace, precursors to a larger and more complex framework that most consider the ultimate promise of the IoT: things connecting, interacting, sharing, storing, and over time perhaps learning and predicting based on habits, behaviors, location, preferences, purchases and more. In his session at @ThingsExpo, Tom Wesselman, Director of Communications Ecosystem Architecture at Plantronics, will examine the still nascent IoT as it is coalescing, including what it is today, what it might ultimately be, the role of wearable tech, and technology gaps stil...
With several hundred implementations of IoT-enabled solutions in the past 12 months alone, this session will focus on experience over the art of the possible. Many can only imagine the most advanced telematics platform ever deployed, supporting millions of customers, producing tens of thousands events or GBs per trip, and hundreds of TBs per month. With the ability to support a billion sensor events per second, over 30PB of warm data for analytics, and hundreds of PBs for an data analytics archive, in his session at @ThingsExpo, Jim Kaskade, Vice President and General Manager, Big Data & Ana...