Showing posts with label Observer. Show all posts
Showing posts with label Observer. Show all posts

Friday, June 26, 2015

Reactive Programming: Reactive Extensions Library (Starting to learn this reactive thing)



I found a great video explanation of Reactive Programming in the flavour of Reactive Extensions Library done by Jafar Husain (Technical Lead at Netflix) . Here's a brief summary of the main concepts I extracted.

Event Reacting
The term Reactive can be associated by analogy with someone throwing a lot of balls at you at the same time. You need to react quickly to try to catch them all. You don't have control over when you want the balls to be thrown at you, you just need to be prepared to catch them.

To start understanding Rx (Reactive Extensions), we need to have two basic design patterns in consideration: Iterator and Observer.

In Java we have an Iterable interface, which is an interface any type of collection can implement to allow consumers of a collection obtain items one at the time. That's how in Java we can use the foreach operator to traverse a collection. Because the interface provides a specific Iterator for the data type we are consuming in the collection. Three things can happen when using an iterator:
  1. Get the next item
  2. No more items to consume (end of the data stream)
  3. An error can happen (using exception throwing in languages like Java)
The Observer pattern is another well known design pattern where suscribers get suscribed to subjects in order to be notified when a change occurs in the subject. If we analize this pattern, we can find that it's very similar to the Iterator pattern in the sense that there is a producer and a consumer. Main difference is that in the Observer the producer is in control for when sending the data, while in the Iterator the consumer is in control. He decides when to pull the data from the producer.

Iteratot Observer Fusion
But there are two main things "missing" from the Observer pattern that are present in the Iterator:
  • A way to indicate there is no more data
  • A way to indicate an error ocurred
In the Observer you can only suscribe callback to receive data, but you cannot register callbacks for the event where is indicated no more data is going to be pushed (completion event), or to indicate an error happened. Reactive extensions is basically about unifying the observable pipe with the iterable pipe producing a new type called the Observable. It gives the same semantics of both the Iterable and the Observable.

There are a lot of interesting operations that can be done over iterables. Similar to SQL where there are a lot of operations that can be done over sets: filter, select, order, etc. What about if all the same things we can do over data sets residing in a table, could also be done over events (data arriving in). That is not a dream with Rx. It is possible to write SQL style queries over events. The difference is that the query is evaluated as data arrives. You can evaluate data in real time.

More and more code is becoming evented. We have a lot of asynchrouns calls both in client side (JS), and server side (Node.JS). That adds a lot of complexity in the application trying to handle all these
callbacks. The reactive extensions library provides this new Observable data type that establishes a powerful way to model events. By completing the missing semantics (end of data stream and error), you can now apply operations familiar in collections like map, filter, reduce, merge, zip, etc. All those things that could be done over data streams we can pull, now can be done over streams of data we can push.

Imagine the new Java 8 Stream API over data that is arriving dynamically over events pushed to you. Lot of things can improve in an application like serving faster data to consumers just to cite an example.

Still need more learning for this new paradimg, but at least I think this video covers fundamental concepts we need to have before getting our hands into the code.

Friday, November 2, 2012

Using Observer Pattern to track progress while loading a page

Have you ever been in a site where there is a heavy process that takes a long time in finishing? If the web page is not user friendly designed, you may end it up with an annoying forever loading page. If we want to avoid this feeling of slowness in our pages, we should consider adding a progress indicator in the page to show how much is left until process is finished. To accomplish this we can take advantage of the Observer pattern. To do this we need to run the process asynchronously, or in other words, running it as a different thread. The next diagram shows how the long process is contained in Thread class.

 

The following class is going to emulate a long process by taking various naps.

package com.gsolano.longprocess

import java.util.Observable;

/**
 * Class with an observable mock long progress.
 * @author gsolano
 *
 */
public class LongProcess extends Observable {
 
 /**
  * Keeps the progress of the process.
  */
 protected Float progress;
 
 /**
  * Simulates a long process.
  */
 public void start() {
  int n =10;
  for (int i=0;i <= n; i++) {
   progress = (float)i/(float)n * 100; // Calculates progress.
   try {
    Thread.currentThread();
    Thread.sleep(2000);  
    this.setChanged();
    this.notifyObservers(progress);
   } catch (InterruptedException e) {   
    e.printStackTrace();
   }
  }  
 }
}


The progress of this class is calculated in every iteration, notifying also the observers with the change in the progress. The next class will observe the LongProcess class.

package com.gsolano.longprocess;

import java.util.Observable;
import java.util.Observer;

/**
 * Observer class
 * 
 * @author gsolano
 */
public class LongProcessObserver implements Observer{

 protected Float progress;
 /**
  * Tracks the progress of the long process.
  * @return
  */
 public Float getProgress() {
  return progress;
 }

 public void update(Observable o, Object arg) {  
  progress = (Float) arg;  
 }
}


To complete the diagram shown before, we need to create a class extending from Thread to wrap the LongProcess and be able to launch in a separate thread.

/**
 * 
 * Class to run a LongProcess in a separate thread.
 * 
 * @author gsolano
 *
 */
public class LongProcessThread extends Thread {
 
 private LongProcess longProcess;
 
 public LongProcess getLongProcess() {
  return longProcess;
 }

 public void setLongProcess(LongProcess longProcess) {
  this.longProcess = longProcess;
 }

 @Override
 public void run() {
  if(longProcess != null) {
   longProcess.start();
  }
 }
}


Now, let’s jump to the web application side. In the next struts action class we handle two events:

1.Start the long process:
  a .Long process is created.
  b. Observer is added to the long process.
  c. Long process is run in a separate thread.
  d. Observer is saved in session variable.

2.Send an update on the progress of the long process
  a. Observer is retrieved from session.
  b. Progress value is taken from observer and written to response.

import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;
import org.apache.struts.action.Action;
import org.apache.struts.action.ActionForm;
import org.apache.struts.action.ActionForward;
import org.apache.struts.action.ActionMapping;

public class FooProgressAction extends Action{
 
 @Override
 public ActionForward execute(ActionMapping mapping, ActionForm form,
   HttpServletRequest request, HttpServletResponse response)
   throws Exception {
  
  String action = request.getParameter("action");
  
  if(action != null) {
   if(action.equalsIgnoreCase("progress")) { // If action is ajax request to get progress.
    // Get the observer.
    LongProcessObserver longProcessObserver = (LongProcessObserver)
      request.getSession().getAttribute("observer");
    if(longProcessObserver != null) {
     // Get the progress from the observer.
     Float progress = longProcessObserver.getProgress();
     if(progress != null) {
      // Send the progress to the page.
      response.getWriter().write(progress.toString());
     }
     return null;
    }
   } else if(action.equalsIgnoreCase("start")) { // Did someone click the start button?
    launchLongProcess(request); 
   }   
  }
  return mapping.findForward("success"); 
 }

 private void launchLongProcess(HttpServletRequest request) {
  LongProcess longProcess = new LongProcess();
  LongProcessObserver observer = new LongProcessObserver();
  // Add the observer to the long process.
  longProcess.addObserver(observer);
  // Launch long process in a thread.
  LongProcessThread longProcessThread = new LongProcessThread();
  longProcessThread.setLongProcess(longProcess);
  longProcessThread.start();  
  // Keep the observer in session.
  request.getSession().setAttribute("observer", observer);
  // Send a flag indicating that party just started!
  request.setAttribute("processStarted", true);
 }
}

In the client side we just need some logic to start the Ajax cycle to ask for progress update until it reaches the 100%.

<%@ taglib uri="/WEB-INF/tld/c.tld" prefix="c" %>

<html>
<head>
 <script language="Javascript">
 var seconds = 1;  
 var run = false;
 var ajaxURL;
 
 function checkProgress(url) {
  if(typeof url != 'undefined') {
   ajaxURL = url;
  } 
    
  var xmlHttp;
  try {
   xmlHttp = new XMLHttpRequest(); // Firefox, Opera 8.0+, Safari
  } catch (e) {
   try {
    xmlHttp = new ActiveXObject("Msxml2.XMLHTTP"); // Internet Explorer
   } catch (e) {
    try {
     xmlHttp = new ActiveXObject("Microsoft.XMLHTTP");
    } catch (e) {
     alert("Ajax not supported");
     return false;
    }
   }
  } 
  xmlHttp.onreadystatechange = function() {
   if (xmlHttp.readyState == 4 ) {
    var progress = xmlHttp.responseText;
    if(progress == 100.0) {
     document.getElementById('progress').innerHTML = "Finished!!"; 
     return;
    }else {
     if(progress) {
      document.getElementById('progress').innerHTML = progress + "%";
     }    
     setTimeout('checkProgress()', seconds * 1000);
    }
   }
  };
  xmlHttp.open("GET", ajaxURL, true);
  xmlHttp.send(null);
 }
 </script>
</head>
 <body>
 <div style="position: absolute; left:40%; text-align:center; border: 1px solid; margin: 20px; padding:20px; width: 150px;">
  <form action="${pageContext.request.contextPath}/longProcess.do">
   <input type="hidden" name="action" value="start" />
   <input type="submit" value="Start!" />
  </form>
  
  <div id="progress"></div>
  
  <c:if test="${not empty processStarted}">
   <script language="Javascript">
    setTimeout('checkProgress(\'${pageContext.request.contextPath}/longProcess.do?action=progress\')', 1000);
   </script>
  </c:if>
 </div>
 </body>
</html>

Result: