Welcome!

Java IoT Authors: Liz McMillan, Pat Romanski, Dana Gardner, Harry Trott, Elizabeth White

Related Topics: Java IoT

Java IoT: 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
"We work in the area of Big Data analytics and Big Data analytics is a very crowded space - you have Hadoop, ETL, warehousing, visualization and there's a lot of effort trying to get these tools to talk to each other," explained Mukund Deshpande, head of the Analytics practice at Accelerite, in this SYS-CON.tv interview at 18th Cloud Expo, held June 7-9, 2016, at the Javits Center in New York City, NY.
Cloud Expo, Inc. has announced today that Andi Mann returns to 'DevOps at Cloud Expo 2016' as Conference Chair The @DevOpsSummit at Cloud Expo will take place on November 1-3, 2016, at the Santa Clara Convention Center in Santa Clara, CA. "DevOps is set to be one of the most profound disruptions to hit IT in decades," said Andi Mann. "It is a natural extension of cloud computing, and I have seen both firsthand and in independent research the fantastic results DevOps delivers. So I am excited t...
Basho Technologies has announced the latest release of Basho Riak TS, version 1.3. Riak TS is an enterprise-grade NoSQL database optimized for Internet of Things (IoT). The open source version enables developers to download the software for free and use it in production as well as make contributions to the code and develop applications around Riak TS. Enhancements to Riak TS make it quick, easy and cost-effective to spin up an instance to test new ideas and build IoT applications. In addition to...
IoT is rapidly changing the way enterprises are using data to improve business decision-making. In order to derive business value, organizations must unlock insights from the data gathered and then act on these. In their session at @ThingsExpo, Eric Hoffman, Vice President at EastBanc Technologies, and Peter Shashkin, Head of Development Department at EastBanc Technologies, discussed how one organization leveraged IoT, cloud technology and data analysis to improve customer experiences and effi...
Internet of @ThingsExpo has announced today that Chris Matthieu has been named tech chair of Internet of @ThingsExpo 2016 Silicon Valley. The 6thInternet of @ThingsExpo will take place on November 1–3, 2016, at the Santa Clara Convention Center in Santa Clara, CA.
Presidio has received the 2015 EMC Partner Services Quality Award from EMC Corporation for achieving outstanding service excellence and customer satisfaction as measured by the EMC Partner Services Quality (PSQ) program. Presidio was also honored as the 2015 EMC Americas Marketing Excellence Partner of the Year and 2015 Mid-Market East Partner of the Year. The EMC PSQ program is a project-specific survey program designed for partners with Service Partner designations to solicit customer feedbac...
The cloud promises new levels of agility and cost-savings for Big Data, data warehousing and analytics. But it’s challenging to understand all the options – from IaaS and PaaS to newer services like HaaS (Hadoop as a Service) and BDaaS (Big Data as a Service). In her session at @BigDataExpo at @ThingsExpo, Hannah Smalltree, a director at Cazena, provided an educational overview of emerging “as-a-service” options for Big Data in the cloud. This is critical background for IT and data profession...
"There's a growing demand from users for things to be faster. When you think about all the transactions or interactions users will have with your product and everything that is between those transactions and interactions - what drives us at Catchpoint Systems is the idea to measure that and to analyze it," explained Leo Vasiliou, Director of Web Performance Engineering at Catchpoint Systems, in this SYS-CON.tv interview at 18th Cloud Expo, held June 7-9, 2016, at the Javits Center in New York Ci...
Ask someone to architect an Internet of Things (IoT) solution and you are guaranteed to see a reference to the cloud. This would lead you to believe that IoT requires the cloud to exist. However, there are many IoT use cases where the cloud is not feasible or desirable. In his session at @ThingsExpo, Dave McCarthy, Director of Products at Bsquare Corporation, will discuss the strategies that exist to extend intelligence directly to IoT devices and sensors, freeing them from the constraints of ...
Connected devices and the industrial internet are growing exponentially every year with Cisco expecting 50 billion devices to be in operation by 2020. In this period of growth, location-based insights are becoming invaluable to many businesses as they adopt new connected technologies. Knowing when and where these devices connect from is critical for a number of scenarios in supply chain management, disaster management, emergency response, M2M, location marketing and more. In his session at @Th...
Extracting business value from Internet of Things (IoT) data doesn’t happen overnight. There are several requirements that must be satisfied, including IoT device enablement, data analysis, real-time detection of complex events and automated orchestration of actions. Unfortunately, too many companies fall short in achieving their business goals by implementing incomplete solutions or not focusing on tangible use cases. In his general session at @ThingsExpo, Dave McCarthy, Director of Products...
The Internet of Things will challenge the status quo of how IT and development organizations operate. Or will it? Certainly the fog layer of IoT requires special insights about data ontology, security and transactional integrity. But the developmental challenges are the same: People, Process and Platform and how we integrate our thinking to solve complicated problems. In his session at 19th Cloud Expo, Craig Sproule, CEO of Metavine, will demonstrate how to move beyond today's coding paradigm ...
There are several IoTs: the Industrial Internet, Consumer Wearables, Wearables and Healthcare, Supply Chains, and the movement toward Smart Grids, Cities, Regions, and Nations. There are competing communications standards every step of the way, a bewildering array of sensors and devices, and an entire world of competing data analytics platforms. To some this appears to be chaos. In this power panel at @ThingsExpo, moderated by Conference Chair Roger Strukhoff, Bradley Holt, Developer Advocate a...
Apixio Inc. has raised $19.3 million in Series D venture capital funding led by SSM Partners with participation from First Analysis, Bain Capital Ventures and Apixio’s largest angel investor. Apixio will dedicate the proceeds toward advancing and scaling products powered by its cognitive computing platform, further enabling insights for optimal patient care. The Series D funding comes as Apixio experiences strong momentum and increasing demand for its HCC Profiler solution, which mines unstruc...
SYS-CON Events has announced today that Roger Strukhoff has been named conference chair of Cloud Expo and @ThingsExpo 2016 Silicon Valley. The 19th Cloud Expo and 6th @ThingsExpo will take place on November 1-3, 2016, at the Santa Clara Convention Center in Santa Clara, CA. "The Internet of Things brings trillions of dollars of opportunity to developers and enterprise IT, no matter how you measure it," stated Roger Strukhoff. "More importantly, it leverages the power of devices and the Interne...
In addition to all the benefits, IoT is also bringing new kind of customer experience challenges - cars that unlock themselves, thermostats turning houses into saunas and baby video monitors broadcasting over the internet. This list can only increase because while IoT services should be intuitive and simple to use, the delivery ecosystem is a myriad of potential problems as IoT explodes complexity. So finding a performance issue is like finding the proverbial needle in the haystack.
Machine Learning helps make complex systems more efficient. By applying advanced Machine Learning techniques such as Cognitive Fingerprinting, wind project operators can utilize these tools to learn from collected data, detect regular patterns, and optimize their own operations. In his session at 18th Cloud Expo, Stuart Gillen, Director of Business Development at SparkCognition, discussed how research has demonstrated the value of Machine Learning in delivering next generation analytics to imp...
Whether your IoT service is connecting cars, homes, appliances, wearable, cameras or other devices, one question hangs in the balance – how do you actually make money from this service? The ability to turn your IoT service into profit requires the ability to create a monetization strategy that is flexible, scalable and working for you in real-time. It must be a transparent, smoothly implemented strategy that all stakeholders – from customers to the board – will be able to understand and comprehe...
The cloud market growth today is largely in public clouds. While there is a lot of spend in IT departments in virtualization, these aren’t yet translating into a true “cloud” experience within the enterprise. What is stopping the growth of the “private cloud” market? In his general session at 18th Cloud Expo, Nara Rajagopalan, CEO of Accelerite, explored the challenges in deploying, managing, and getting adoption for a private cloud within an enterprise. What are the key differences between wh...
The IoT is changing the way enterprises conduct business. In his session at @ThingsExpo, Eric Hoffman, Vice President at EastBanc Technologies, discussed how businesses can gain an edge over competitors by empowering consumers to take control through IoT. He cited examples such as a Washington, D.C.-based sports club that leveraged IoT and the cloud to develop a comprehensive booking system. He also highlighted how IoT can revitalize and restore outdated business models, making them profitable ...