Saturday, November 28, 2015

Java 8 - Streams (Part - I)

Stream represents a sequence of elements which can calculate or compute on-demand. Stream is an interface like an Iterator but it can do parallel processing. In other words, we can say Stream is a lazy collection where the values are computed on-demand.

How to Create

Streams are defined in java.util.stream package. Stream can be obtained using various options but a simple way is to get Stream is from Collection Interface. A default method defined as stream in Collection Interface to get it from the respective collection built. This is one of the best example to explain why the default methods are introduced in Java 8. See my previous post for more details on default methods. Collection Interface has two default methods for streams
  • stream(): Returns a stream associated with the collection for processing
  • parallelStream(): Same as stream method but returns for parallel processing. 
Stream here is the generic type. There are few streams defined for the primitive types as IntStream, DoubleStream, LongStream etc.

How to Use

As mentioned earlier - Streams are like Iterators or more than Iterators. Streams can be used to find an element, get first element, sort the elements etc. Apart from using streams on Collections, they can be used to operate on Paths, files, range of numbers etc. We will see few of them with examples

Streaming Files

BufferedReader class has got a new method lines() to return the lines in that file as Stream. See below
        try (FileReader fr = new FileReader("/tmp/sample.txt");
                BufferedReader br = new BufferedReader(fr))
        {
            br.lines().forEach(System.out::println);
        }
The above code returns the Stream which represents the sequence of lines in the file /tmp/sample.txt and, prints them. To add to it, Stream has got a method forEach to iterate over it. The Files class has also has a method lines() to read the Path as stream.
        try (Stream<String> st = Files.lines(Paths.get("/tmp/sample.txt")))
        {
            st.forEach(System.out::println);
        }

Streaming Patterns

Pattern can also return a Stream object with the matched values. Lets see that
        Pattern p = Pattern.compile("-");
        p.splitAsStream("5-13-93").forEach(System.out::println);
This will split the pattern 5-13-93 into three integers (5,13 and 93) and prints them.

Streaming Range

Stream interface has method to find range of primitive types (Int, Double or Long) as a Stream itself. We can use the respective Stream Interface (IntStream, DoubleStream and LongStream) to get the range associated with it. See - for example
IntStream.range(10, 20).forEach(System.out::println);

Stream Functions

Stream provides various methods to operate on Collections. We will discuss few of them here with examples. Most of the methods of Stream interface take lambda expression as argument (i.e. a method name - See Lambda Expressions for more details).

of:

This is a static method defined in Stream Interface to create a generic Stream. The method exists in other specific Streams like Double, Long etc to create its own type.
Stream<String> stream = Stream.of("Veeresh", "Blog", "Stream")

forEach:

The function will iterate through the Stream. All the above example has got forEach with prints to console

sort:

This method sorts the Stream elements associated with it in natural order (i.e. ascending order). See here
Stream.of("Veeresh", "Blog", "Stream").sorted().forEach(System.out::println);
The above code prints the sorted names. sort() is a very lazy function, means it doesn't effect until you call any other method after/before sort. In the above example - When the sort method was executed nothing changed on the Stream. Once the foreach method was called - the sorting was done.

Apart from these - stream has few other methods which are more powerful and useful which reduces developer effort, Increases readability and reduces LOC. We will discuss those in the next post. For now - this is it.

Happy Learning!!!!

Sunday, October 11, 2015

Default Methods in Java

Default Method is one of the features of Java 8. Default Method is
  • A method - which has a keyword default before it
  • Has to be defined only in Interface
  • Must be implementation. 
  • Can be overridden in Implemented Classes but has to be defined without default keyword.
  • An Interface can have multiple default methods
Default method may break the contract of an Interface because technically interface should consists of all abstract methods and no implementations. The main reason behind the concept of default method is to add steam method to collections (especially to interface) without breaking the implementation classes.

Example

  • As mentioned earlier, default method is a method with implementation in an Interface
public interface AInterface 
{
    public default void print()
    {
        System.out.println("Default Method in Interface A");
    }
}
  • If a class implements it, then it may or may not override the default method along with it's own methods.
public class AImplementation implements AInterface
{
    public void otherMethod()
    {
        System.out.println("New Method Defined in AImplementation");
    }
}
  • Calling a default is same as like other methods as below
AImplementation a = new AImplementation();
a.print();
a.otherMethod();
  • As multiple interfaces can be implemented by one class, but there is a probability of same default method with the same name, defined in both interfaces. Let's say another interface with a default method print like below
public interface BInterface
{
    default void print()
    {
        System.out.println("Default Method in Interface B"); 
    }
}
  • Now, if a class implements both AInterface and BInterface, which has same default method print, then compiler throws an error saying "Duplicate default methods named print with the parameters () and () are inherited from the types BInterface and AInterface". So, print method has to be overridden in the implemented Class like below
public class CImplementation implements AInterface, BInterface
{
    @Override
    public void print()
    {
        System.out.println("Default Method in Interface A");
    }
}

More details on the default methods and its use in Java will be explained in the next post.

Happy Learning!!!!

Sunday, September 27, 2015

Apache Camel-JPA

Camel JPA allows to route the Entities as Messages to read and write using Java Persistence Architecture (JPA). Let's quickly look at the configuration required to make Camel to route the messages (in fact entities) using camel-jpa.

Dependency

Add the following dependency along with camel-core to enable Camel JPA
<dependency>
	<groupId>org.apache.camel</groupId>
	<artifactId>camel-jpa</artifactId>
	<version>2.15.3</version>
</dependency> 

Component

Create JPA component using the class org.apache.camel.component.jpa.JpaComponent and inject Entity Manager Factory (and Transaction Manager if required)
<bean id="jpa" class="org.apache.camel.component.jpa.JpaComponent">
	<property name="entityManagerFactory" ref="entityManagerFactory" />
	<property name="transactionManager" ref="transactionManager" />
</bean>

EndPoint

Format of the JPA Endpoint URI to be defined as (Component, Entity Class and Options)
jpa://<fully-classified-entity-class-name?option1=value1&option2=value2.."

Few Notable Parameters 

  • entityType: Default value is Class Name of the entity. The option will override the entity class name.
  • consumeDelete: The table row will be deleted after processing. In order to avoid deleting - set it to false. 
  • consumer.query: Custom query to consume the data from the table. 
  • consumer.namedQuery: Custom named query to consume the data from the table.
  • consumer.nativeQuery: Custom native query to consume the data form the table.
  • consumer.delay: Polling interval in milli-seconds between two polls.

Reading from Endpoint

Selection

By default, JPA will read all the rows in the table. To read specific set of rows in the table - use consumer.query (, consumer.nativeQuery or consumer.namedQuery) to set the selection. 

Preserve Rows

After the row is read from the table, the endpoint deletes the processed row. In order to avoid deleting the row, set the parameter consumeDelete to false. 
If we set this parameter - there is a possibility of reading the same row again because the row still exists on the table. In order to avoid the same row to be picked again - we can update the row by creating a method (in the entity and changing the specific value) by annotating it with @Consumed. 

Sending to Endpoint

Send entity by entity to an endpoint, it writes row by row into the table. We can use aggregator to write a list of entities into the endpoint. In this case, we need to set an extra parameter : entityType as java.util.List to make all the rows persisted to the table. 

Example

Let's read rows from a table and write into the same table by changing few parameters.

Table

mysql> desc camel_test;
+---------+----------+------+-----+---------+-------+
| Field   | Type     | Null | Key | Default | Extra |
+---------+----------+------+-----+---------+-------+
| seq_no  | int(11)  | NO   |     | NULL    |       |
| status  | char(4)  | NO   |     | NULL    |       |
| column1 | char(20) | YES  |     | NULL    |       |
| column2 | char(40) | YES  |     | NULL    |       |
| column3 | int(11)  | YES  |     | NULL    |       |
| column4 | datetime | YES  |     | NULL    |       |
+---------+----------+------+-----+---------+-------+

Configuration

We will see the example using Hibernate and Spring JPA with MYSQL

persitance.xml

<persistence xmlns="http://java.sun.com/xml/ns/persistence" version="1.0">
    <!--Persistence Unit for Mysql database-->
    <persistence-unit name="testMysql" transaction-type="RESOURCE_LOCAL">
        <provider>org.hibernate.ejb.HibernatePersistence</provider>
        <properties>
            <property name="hibernate.dialect" value="org.hibernate.dialect.MySQL5InnoDBDialect"/>
            <property name="hibernate.show_sql" value="true"/>
        </properties>
    </persistence-unit>
</persistence>

Datasource and Entity manager

<bean id="mysqltestDataSource"
	class="org.springframework.jdbc.datasource.DriverManagerDataSource">
	<property name="driverClassName" value="${jdbc.driverClassName}" />
	<property name="url" value="${jdbc.testurl}" />
	<property name="username" value="${jdbc.username}" />
	<property name="password" value="${jdbc.password}" />
</bean>

<context:property-placeholder location="classpath:jdbc.properties" />

<bean id="entityManagerFactory"
	class="org.springframework.orm.jpa.LocalContainerEntityManagerFactoryBean">
	<property name="dataSource" ref="mysqltestDataSource" />
	<property name="persistenceUnitName" value="testMysql" />
</bean>

Camel JPA

<bean id="jpa" class="org.apache.camel.component.jpa.JpaComponent">
	<property name="entityManagerFactory" ref="entityManagerFactory" />
</bean>

Entity Class

package com.test.entity;
//Imports here.. 
@Entity
@Table(name = "camel_test")
@NamedQuery(name = "selectQuery", query = "select m from CamelTest m where m.status = 'P'")
public class CamelTest
{
    @Id
    @Column(name="seq_no")
    private Integer seq;
    
    @Column(name="status")
    private String status;
    
    @Column(name="column1")
    private String column1;
    
    @Column(name="column2")
    private String column2;
    
    @Column(name="column3")
    private Integer column3;
    
    @Column(name="column4")
    private Date column4;

    //Getters and Setters here
    
    @Consumed
    public void afterConsume()
    {
        setStatus("D");
    }
}

Points to be noted

  • The class is a JPA entity with one Named Query which will be used by JPA endpoint to read the data from the Table.
    • It reads the data from table whose status is 'P'
  • Only one extra method added which annotated with @Consumed (This will be executed after the entity is processed.
    •  After Camel JPA processing the record, the entity is changing the status to 'D' to avoid it from further reading

Finally Camel Route

<camel:camelContext id="camelContext">
	<camel:route>
		<camel:from uri="jpa://com.test.entity.CamelTest?consumer.namedQuery=selectQuery&amp;consumer.delay=10000&amp;consumeDelete=false" />
		<camel:to uri="direct:sample" />
	</camel:route>
</camel:camelContext>

Points to be noted

  • jpa is the JPA Component of Camel JPA
  • com.test.entity.CamelTest is the entity class created corresponding to camel_test table (as above)
  • Three options are set
    • consumeDelete to false to make the entity preserved after processing
    • consumer.dealy is set to 10 seconds to poll the table for every 10 seconds
    • consumer.namedQuery is set to read a specific set of rows (see NamedQuery annotation on the Entity). 
The above route just reads the rows from the table using the named query (selectQuery) with a polling delay of 10 secs, maps the entities, executes the @Consumed method (which modifies the status column from P to D so that the next poll will not pick the same row)  in the entity class after processing and finally sends the processed message to direct:sample (one by one)

Read in a batch

To read the rows in a batch, use aggregator

Aggreagor

package com.test.camel.aggregator;
//Imports
public class SampleAggregator implements AggregationStrategy
{
    @SuppressWarnings("unchecked")
    public Exchange aggregate(Exchange oldExchange, Exchange newExchange)
    {
        List<CamelTest> list = null;
        if(oldExchange == null || oldExchange.getIn() == null)
        {
            list = new ArrayList<CamelTest>();
            list.add(newExchange.getIn().getBody(CamelTest.class));
        }
        else
        {
            list = oldExchange.getIn().getBody(List.class);
            list.add(newExchange.getIn().getBody(CamelTest.class));
        }
        newExchange.getIn().setBody(list);
        return newExchange;
    }
}

Modified Route as below with Aggregator

<camel:route>
	<camel:from uri="jpa://com.test.entity.CamelTest?consumer.namedQuery=selectQuery&amp;consumer.delay=10000" />
	<camel:aggregate strategyRef="aggregateStrategy" completionInterval="2000" completionSize="20" >
		<camel:correlationExpression>
			<camel:constant>true</camel:constant>	
		</camel:correlationExpression>
		<camel:to uri="direct:sample" />
	</camel:aggregate>
</camel:route>

Write to Table in Batch

Modify the same route to write a duplicate entry into the table by changing the Sequence Number. See the processor below

Processor

package com.test.camel.processor;
//Imports
public class SampleProcessor implements Processor
{
    public void process(Exchange exchange) throws Exception
    {
        //Processing goes here. Set the output to same as input. (Dummy Processing) 
        List<CamelTest> list = exchange.getIn().getBody(List.class);
        
        for(CamelTest t : list)
        {
            t.setSeq(t.getSeq() + 1000); 
        }
        exchange.getOut().setBody(list);     
    }
}

  • Processor - processor is used here to change the primary key of the entities to write it back (Other we get primary key violation error) - this is just for the example purpose - added 1000 for each sequence number.

Modified Route

<camel:route>
	<camel:from uri="jpa://com.test.entity.CamelTest?consumer.namedQuery=selectQuery&amp;consumer.delay=10000" />
	<camel:aggregate strategyRef="aggregateStrategy" completionInterval="2000" completionSize="20" >
		<camel:correlationExpression>
			<camel:constant>true</camel:constant>	
		</camel:correlationExpression>
		<camel:process ref="processor" />	
		<camel:to uri="jpa://com.test.entity.CamelTest?entityType=java.util.List" />
	</camel:aggregate>
</camel:route>
  • URI on the route (camel:to) has a parameter called entityType as java.util.List to save a list of entities back to table. (See Sending to Endpoint above)
Happy Learning!!!

Sunday, September 13, 2015

Apache Camel - Simple Routing Example

This post is continuation of the previous post Apache Camel overview. In this post, we will see how to define a sample route using Apache Camel (with Spring Integration).

Dependencies

Apache Camel can work easily with Spring. The following dependencies are required to make Apache Camel to work with Spring.
  <dependency>
       <groupId>org.apache.camel</groupId>
       <artifactId>camel-core</artifactId>
       <version>2.15.3</version>
  </dependency>
  <dependency>
       <groupId>org.apache.camel</groupId>
       <artifactId>camel-spring</artifactId>
       <version>2.15.3</version>
  </dependency>
camel-core is the actual dependency and camel-spring is required for the spring integration.

Spring Configuration

To add camel to Spring Configuration, add the following

Namespace

xmlns:camel="http://camel.apache.org/schema/spring"

Schema Location

xsi:schemaLocation="http://camel.apache.org/schema/spring
http://camel.apache.org/schema/spring/camel-spring.xsd"

Combined Spring XML

<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
    xmlns:camel="http://camel.apache.org/schema/spring"
    xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://www.springframework.org/schema/beans 
              http://www.springframework.org/schema/beans/spring-beans.xsd
              http://camel.apache.org/schema/spring 
              http://camel.apache.org/schema/spring/camel-spring.xsd">

</beans>         

Camel Context and Route

Camel context is the camel runtime. In Spring, we need to define the Camel Context to define the Route. Within the context, route must be defined.
Problem Statement: Read a file from a directory, process it and copy the processed file to a destination folder.

Define the Processor

package com.test.camel.processor;

import org.apache.camel.Exchange;
import org.apache.camel.Processor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class SampleProcessor implements Processor
{
    private static final Logger LOG = LoggerFactory.getLogger(SampleProcessor.class);

    public void process(Exchange exchange) throws Exception
    {
        LOG.info("Input is "+exchange.getIn().getBody());
        //Processing goes here. Set the output to same as input. (Dummy Processing) 
        exchange.setOut(exchange.getIn());
    }

}

Points to be noted:

  • Processor is a class which implements org.apache.camel.Processor. 
  • Processor implements a method which takes only one argument (Exchange) - which consists of three messages (in, out and exception) and few other properties related to the route and it's camel context.

Define the Bean

     <bean id="processor" class="com.test.camel.processor.SampleProcessor" />

Define the Route

 <camel:camelContext id="camelContext">
      <camel:route id="file-route">
            <camel:from uri="file:/tmp/input.test" />
            <camel:process ref="processor" />
            <camel:to uri="file:/tmp/output.test" />
      </camel:route>
  </camel:camelContext>

Points to be noted:

  • Route is defined inside the Spring Config (using DSL). It also be defined using RouteBuilder. (We will concentrate only on DSL).
  • Define the route inside the Camel Context 
  • A very basic route consists of From (Where to fetch from), Processor (Which processor the input) and To (Where to write the processed input). 
  • Processor is optional - if we don't want to do any processing. 
  • Here in this route (file-route), files will be read from directory /tmp/input.test, and will be written into /tmp/output.test directory after processing.
  • Camel Context and Route(s) will be automatically started once the spring application context is loaded and started. So, once spring configuration file is loaded, camel looks for the files in the input directory. 
  • Camel Context and Route(s) will be shutdown once the application context is closed.

Exception Handling

Apache Camel provides try, catch and finally (optional) blocks to be defined within a route to handle the exceptions while routing or processing. The above route can be re-defined with try-catch blocks to check for Exception 
 <camel:camelContext id="camelContext">
      <camel:route id="file-route">
            <camel:from uri="file:/tmp/input.test" />
            <camel:doTry>
                 <camel:process ref="processor" />
                 <camel:to uri="file:/tmp/output.test" />
                 <camel:doCatch>
                      <camel:exception>java.lang.Exception</camel:exception>
                      <camel:log message="Exception while processing "></camel:log>
                      <camel:to uri="file:/tmp/error.test" />
                 </camel:doCatch>
            </camel:doTry>
      </camel:route>
  </camel:camelContext>
In the above case, if an exception is occurred while routing/processing, the exception will be caught, a log will be written and the file will be moved to another directory /tmp/output.test.

Full Spring Configuration

Complete spring configuration file is as below
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
      xmlns:camel="http://camel.apache.org/schema/spring" 
      xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
      xsi:schemaLocation="http://www.springframework.org/schema/beans 
           http://www.springframework.org/schema/beans/spring-beans.xsd
           http://camel.apache.org/schema/spring 
           http://camel.apache.org/schema/spring/camel-spring.xsd">

     <bean id="processor" class="com.test.camel.processor.SampleProcessor" />

     <camel:camelContext id="camelContext">
         <camel:route id="file-route">
              <camel:from uri="file:/tmp/input.test" />
              <camel:doTry>
                   <camel:process ref="processor" />
                   <camel:to uri="file:/tmp/output.test" />
                   <camel:doCatch>
                        <camel:exception>java.io.IOException</camel:exception>
                        <camel:log message="Exception while processing "></camel:log>
                        <camel:to uri="file:/tmp/error.test" />
                   </camel:doCatch>
              </camel:doTry>
         </camel:route>
     </camel:camelContext>
</beans>
         
In the next post, we will see more details about routing. Happy Learning!!!!

Sunday, August 30, 2015

Apache Camel and KeyWords

To be very simple, Apache Camel is an open source Java Based API which implements most of the commonly used EIPs. So, the next question is what is an EIP?. Let me briefly explain

Enterprise Integration Patterns

EIPs is a set of design patterns (65 design patterns) on how to integrate different enterprise application running on various technologies and makes them to communicate.  These patterns are explained in a book by Gregor Hohpe and Bobby Woolf. So, EIPs is a book in which, the ways to create communication between applications and integrate them are defined. Check the link for more details.

Apache Camel

Let's come back to Apache Camel. As I mentioned at the beginning, Apache Camel is an implementation of EIPs. So, we can integrate the enterprise applications using Camel. It supports almost all types of protocols (FTP, HTTP, JMS, WebService, etc,), to communicate either by itself or by leveraging it. See below an overview of Camel
Using Camel, we can define a set of route(s) within Camel Context to create the communication. The received messages can be filtered, processed, validated, aggregated, split or routed to different other systems.

Camel Context

Camel Context creates runtime for Apache Camel.

EndPoint

As the name indicates, an endpoint is the final or intermediate source/destination of the communication. The endpoint can be a physical address like FTP server or a logical address like JMS Queue, WebService etc. Apache Camel uses URI (of course, Uniform Resource Identifier) to define endpoints. Ex: queue:SAMPLE for a Queue.

Exchange

Exchange is analogous to Message in JMS Specification. But the Exchange carries three messages i.e. Incoming, Outgoing and an Exception for processing. These three are defined as in, out and fault messages. Each of theses Messages has it's own headers.

Template  

CamelTemplate is a class which reads from/writes to an EndPoint. Earlier versions of the camel has the name as CamelClient, but the convention is changed in later versions to be in sync with the other implementations (like Spring JmsTemplate). Template reads/writes Exchanges to/from EndPoint.

Component

Component is a factory class through which you can create an instance of an EndPoint. JmsComponent is the factory for creating JMS Endpoints using a method called JmsComponent.createEndpoint();

Processor

The points discussed so far are very basic to integrate at-least two applications as it is (exchange the messages as it is with no processing). In order to process the message before sending to destination application, we need processor. The processor is an interface which has only one method called void process(Exchange exchange);

Route

A route is a step by step movement of the message between the two endpoints (including exception handling). There are two ways to define a route in Apache Camel. One is using the XML file (like spring bean file). Second is by using Java DSL (Domain Specific Language).

In the next post, we will see how to use these concepts to create a route and process the message between two endpoints (within the Camel Context).

Happy Learning!!!!!