Monday, May 20, 2019

A Basic Authorized Service in Amazon AWS

In this post, I will produce one very unimpressive web-based application: a login page giving access to a second page that displays the current date and time. I will use the following AWS technologies to accomplish this:

  • The Simple Storage Service that will serve as a web server hosting publicly available html and javascript pages.
  • Lambda Functions that will provide the login operation, the current date and the authorizer, which serves to prevent unauthorized access to the current date.
  • API Gateway, which will expose the service offered by the Lambda functions as a REST interface.

Amazon documentation is excellent in general and we will follow it whenever possible. However, it misses more elaborate combinations of services like this. The following diagram illustrates the interactions we are trying to accomplish:



The browser gets the client application from an S3 bucket (step 1). This application includes an HTML file and a React application. This latter application uses JavaScript to do the login (step 2) and get the date and time from an online service (step 3). Since this service is protected by an authorization token, this involves the verification of the token (step 4), before the actual access to the date service (step 5). A major problem with this idea is that the services and the pages are hosted in different domains, which makes the browser unable to look for the services in JavaScript, due to Cross-Origin Resource Sharing (CORS) restrictions. We will need to take that into consideration when designing the services and configuring the API Gateway. We may also put the S3 bucket behind the API Gateway, to serve the entire content from the same domain. However, that configuration is slightly more complicated than the alternative I present here.

Let us start with the client application. It involves two files. I called hello.html to the first, and counter.js (inside a js directory) to the other. The HTML file is as follows. Note that it refers the other file:

<!DOCTYPE html>
<html>
    <head>
        <title>
            Hello Authenticated World!
        </title>
        <script src="https://unpkg.com/react@16/umd/react.development.js" crossorigin></script>
  <script src="https://unpkg.com/react-dom@16/umd/react-dom.development.js" crossorigin></script>
  <script src="https://unpkg.com/babel-standalone@6/babel.min.js" crossorigin></script>
    <script type="text/babel" src="js/counter.js"></script>
    </head>
    <body>
        <h1>Hello Authenticated World!</h1>
        <div id="counter"/>
    </body>
</html>

The HTML code has a div element that is replaced by a React component in the counter.js file:

class Login extends React.Component {
    constructor(props) {
        super(props);
        this.state = {login : '', password : ''}

        this.handleLoginChange = this.changeLogin.bind(this);
        this.handlePasswordChange = this.changePassword.bind(this);
    }

    changeLogin(event) {
        this.setState({ login : event.target.value })
    }

    changePassword(event) {
        this.setState({ password : event.target.value })
    }

    render () {
        return (
            <form>
                Login <input type="text" value={this.state.login} onChange={this.handleLoginChange}/> <br/>
                Password <input type="password" value={this.state.password} onChange={this.handlePasswordChange}/> <br/>
                <input type="button" value="submit" onClick={this.props.login.bind(this.props.parent, this.state.login, this.state.password)}/>
            </form>
        )
    }
}


class Watch extends React.Component {
    constructor(props) {
        super(props);
        this.state = {date : ''};
    }

    componentWillMount() {
        var theobject = this
        console.log('this =', this)
        console.log('this.props =', this.props)
  fetch(this.props.timeserver, {
            headers:{
              'authorizationToken' : JSON.stringify({ 'token' : this.props.token })
            }
        }).then(data => data.json())
          .then(thedate => theobject.setState({date : thedate}))
    .catch(function(error) {
            console.log('There has been a problem with your fetch operation: ', error.message);
           });
    }

    render() {
        console.log('date =', this.state.date)
        return <h2>{this.state.date}</h2>
    }
}


class Page extends React.Component {
    constructor(props) {
        super(props);
        this.state = {token : undefined, errormessage : undefined };
    }

    authenticated(token) {
        this.setState({ token : token, errormessage : undefined })
    }

    failedauthenticated() {
        this.setState({ token : undefined, errormessage : 'Authentication Error' })
    }

    doLogin(login, password) {
        var theobject = this
        var formparameters = {
            method: 'POST', // or 'PUT'
            body: JSON.stringify({'login' : login, 'password' : password}),
            headers:{
              'Content-Type': 'application/json'
            }
        }
  fetch(this.props.loginserver, formparameters).then(function(data) {
            if(data.status!==200) {
                theobject.failedauthenticated()
                throw new Error(data.status)
            }
            else {
                var json = data.json();
                return json;
            }
  }).then(function(thetoken) {
            console.log('message =', thetoken)
            if ('token' in thetoken)
                theobject.authenticated(thetoken['token'])
  }).catch(function(error) {
            console.log('There has been a problem with your fetch operation: ', error.message);
        });
    }

    render() {
        console.log(this.state.token)
        if (this.state.errormessage != undefined)
            var errormessage = <h2>{this.state.errormessage}</h2>
        else
            var errormessage = <div/>
        if (this.state.token == undefined)
            return (
                <div>
                    {errormessage}
                    <Login parent={this} login={this.doLogin}/>
                </div>
            )
        else
            return (
                <div>
                    {errormessage}
                    <Watch timeserver={this.props.timeserver} token={this.state.token}/>
                </div>
            )
    }
}

ReactDOM.render(<Page loginserver="https://your-service.execute-api.eu-west-1.amazonaws.com/production/login" timeserver="https://your-service.execute-api.eu-west-1.amazonaws.com/production/time"/>, counter);


VERY IMPORTANT: The URLs in the end are not valid. You must replace them with the addresses that you will create on API Gateway, below in this document. (Later edit) A student of mine was complaining about the "this.state.date" and claiming that this could be solved with "this.state.date.body". I haven't tried.

You may upload these two files to an S3 bucket. To learn how to do this, please refer to AWS documentation. The crucial part has to do with making these two files available for public access. Please refer to this URL here, or find it looking for "host web site s3" or something like that in your favorite search engine. I include here the policy that I used in my bucket:

{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "PublicReadGetObject",
            "Effect": "Allow",
            "Principal": "*",
            "Action": "s3:GetObject",
            "Resource": "arn:aws:s3:::bucket-name/*"
        }
    ]
}

Essentially, this gives authorization to any principal (user) to access objects in your bucket. Once you succeed, you should be able to point your browser to your new web site to see this:

This page is not ready to do anything yet, we need to do something with the submit button. Let us start with the login lambda function, in Python:

#import jwt
import json


#password = 'Very long, very difficult indeed! It would take ages to attack this very long string. Now the numbers that were missing so far: 34719108435'

def lambda_handler(event, context):
    print(event)
    response = {
        'statusCode': 401
    }
    if 'body' in event:
        print(event['body'])
        contents = json.loads(event['body'])
        if 'login' in contents:
            login = contents['login']
        else:
            login = None
        if 'password' in contents:
            password = contents['password']
        else:
            password = None
        if login != None and login != '' and login == password:
            #code = jwt.encode({'some': 'payload'}, password, algorithm='HS256')
            response = {
                'statusCode': 200,
                #'body': json.dumps({'token' : code.decode('UTF-8')})
                'body': json.dumps({'token' : 'authorized'})
            }

    response['headers'] = {
        'Access-Control-Allow-Origin' : '*', # Required for CORS support to work
        'Access-Control-Allow-Credentials' : True, # Required for cookies, authorization headers with HTTPS 
    }

    return response

See this on how to "create a Lambda function" or search for this exact expression on your search engine.

A few details about this function: to make everything simpler, I am just checking if the login is the same as the password. Needless to say you should never ever implement anything like that. I should have used a standard technology like JSON web tokens (JWT), but that would complicate the exercise a bit. We would need to import a library, following steps like these. Nevertheless, I commented some code that could help you with JWT.

Note also the headers in the response, which enable resource sharing with any origin. This overcomes the different domains of the S3 web site and API Gateway. Without these headers the example will not work.

Once the lambda function is ready, one should go the the API Gateway and link a resource to this lambda function. Check this documentation from AWS for that purpose. Essentially, you need to

  • Create the API.
  • Create the /login resource with CORS enabled.
  • Create a POST method in the /login resource.
  • Deploy the API.
The following figures summarize these steps:






Before actually deploying the service, we may test the /login POST method, by passing the following data in the request body:

{
    "login": "Paul",
    "password": "Paul"
}

The answer should be:

{
  "token": "authorized"
}

and the response headers:

{"Access-Control-Allow-Origin":"*","Access-Control-Allow-Credentials":"true","X-Amzn-Trace-Id":"Root=1-5ce1ad88-4eaf7477d48dd8a4b8ad10f8;Sampled=0"}

Once one deploys the service, the service should enable login, but we still miss the get date function. For that Lambda function, we need the following code (note the headers to enable CORS again):

import json
from datetime import datetime


def lambda_handler(event, context):
    now = datetime.now() # current date and time
    date_time = now.strftime("%m/%d/%Y, %H:%M:%S")
    return {
        'statusCode': 200,
        'body': json.dumps(date_time),
        'headers': {
                    'Access-Control-Allow-Origin' : '*', # Required for CORS support to work
                    'Access-Control-Allow-Credentials' : True, # Required for cookies, authorization headers with HTTPS 
        }
    }

We will associate this Lambda function to a resource called /time in the API Gateway. We enable CORS again:




This time, we need an extra action, to let the 'authorizationToken' go through the /time service. We need to manually enable CORS (checking the box is not enough) and add the 'authorizationToken' to the Access-Control-Allow-Headers list. In the end, we should get something like this in the response:


Don't ever forget to deploy after each change or you will not be able to see anything. In the deploy step you will get the URL to replace the loginserver and timeserver properties in the counter.js file.

One final step is missing: adding an authorizer to protect the /time service. Indeed, without the authorizer you can put its URL on a browser and see the result of the service. We will prevent that from happening, with the help of our final Lambda function:

#import jwt
import json

#password = 'Very long, very difficult indeed! It would take ages to attack this very long string. Now the numbers that were missing so far: 34719108435'


def lambda_handler(event, context):
    if 'authorizationToken' in event:
        #result = jwt.decode(event['authorizationToken'], password, algorithms = ['HS256'])
        result = json.loads(event['authorizationToken'])
        if 'token' in result and result['token'] == 'authorized': #should check time instead
            return generatePolicy('user', 'Allow', event['methodArn'])
        else:
            return generatePolicy('user', 'Deny', event['methodArn'])
    else:
        return 'Unauthorized'



# Help function to generate an IAM policy
def generatePolicy(principalId, effect, resource):
    authResponse = {}
    
    authResponse['principalId'] = principalId
    if effect and resource:
        policyDocument = {}
        policyDocument['Version'] = '2012-10-17'
        policyDocument['Statement'] = []
        statementOne = {}
        statementOne['Action'] = 'execute-api:Invoke'
        statementOne['Effect'] = effect
        statementOne['Resource'] = resource
        policyDocument['Statement'] = [statementOne]
        authResponse['policyDocument'] = policyDocument
    
    # Optional output with custom properties of the String, Number or Boolean type.
    # authResponse['context'] = {
    #     "stringKey": "stringval",
    #     "numberKey": 123,
    #     "booleanKey": True
    # };
    return authResponse

To create a new authorizer use the following parameters. Note the "token source" field, which must match the name of the header we are using to pass the authorization token: authorizationToken.


Then, we must go to the GET method of the /time resource and protect it with the authorizer (you may have to reload the page to see the Token authorizer):


Once you do this, the /time URL is no longer accessible on a browser and we are all set. Once you login with a username equal to the password, you should see this:


Thursday, November 22, 2018

A Standalone REST server in Java

Doing a standalone REST server, i.e., a REST server that runs outside WildFly or Tomcat or any other container is not too difficult in Java. The part that I had trouble with was figuring out the right dependencies in Maven. I will go straight to the code, as I think this is self-explanatory. I have a class that provides the "students" service, which basically keeps a list of students' names. This is basic to say the least, because I keep the list in a static property, as the server may have several threads. Real implementations would probably be backed by some sort of database seen by all threads. But for this example this suffices...

package is.project3.rest;

import java.util.ArrayList;
import java.util.List;

import javax.ws.rs.GET;
import javax.ws.rs.POST;
import javax.ws.rs.Path;
import javax.ws.rs.Produces;
import javax.ws.rs.core.MediaType;

@Path("/students")
public class StudentsKeeper {
 //XXX: we have several threads...
 private static List<String> students = new ArrayList<>();
 
 @POST
 public void addStudent(String student) {
  System.out.println("Called post addStudent with parameter: " + student);
  System.out.println("Thread = " + Thread.currentThread().getName());
  students.add(student);
 }

 @GET
 @Produces(MediaType.APPLICATION_JSON)
 public List<String> getAllStudents() {
  System.out.println("Called getAllStudents");
  System.out.println("Thread = " + Thread.currentThread().getName());
  return students;
 }

 @Path("xpto")
 @GET
 @Produces(MediaType.APPLICATION_JSON)
 public String getAllStudents2() {
  System.out.println("Called getAllStudents 2");
  System.out.println("Thread = " + Thread.currentThread().getName());
  return "test";
 }

}

The server is surprisingly simple:

package is.project3.rest;


import java.net.URI;

import javax.ws.rs.core.UriBuilder;

import org.glassfish.jersey.jdkhttp.JdkHttpServerFactory;
import org.glassfish.jersey.server.ResourceConfig;

public class MyRESTServer {

 private final static int port = 9998;
 private final static String host="http://localhost/";

  public static void main(String[] args) {
  URI baseUri = UriBuilder.fromUri(host).port(port).build();
  ResourceConfig config = new ResourceConfig(StudentsKeeper.class);
  JdkHttpServerFactory.createHttpServer(baseUri, config);
 }
}


You just need to make it run and it will immediately be serving a REST web service on port 9998. More about this ahead. The troubling part was to get Maven right, considering that it had to transparently convert a List of Strings to JSON.

<project xmlns="http://maven.apache.org/POM/4.0.0"
 xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
 xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
 <modelVersion>4.0.0</modelVersion>

 <groupId>is</groupId>
 <artifactId>project3</artifactId>
 <version>0.0.1-SNAPSHOT</version>
 <packaging>jar</packaging>

 <name>project3</name>
 <url>http://maven.apache.org</url>

 <properties>
  <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
  <java.version>10</java.version>
 </properties>

 <dependencies>
  <!-- https://mvnrepository.com/artifact/junit/junit -->
  <dependency>
   <groupId>junit</groupId>
   <artifactId>junit</artifactId>
   <version>4.12</version>
   <scope>test</scope>
  </dependency>

  <!-- https://mvnrepository.com/artifact/javax.xml.bind/jaxb-api -->
  <dependency>
   <groupId>javax.xml.bind</groupId>
   <artifactId>jaxb-api</artifactId>
   <version>2.3.0</version>
  </dependency>
  <!-- https://mvnrepository.com/artifact/com.sun.xml.bind/jaxb-core -->
  <dependency>
   <groupId>com.sun.xml.bind</groupId>
   <artifactId>jaxb-core</artifactId>
   <version>2.3.0.1</version>
  </dependency>
  <dependency>
   <groupId>com.sun.xml.bind</groupId>
   <artifactId>jaxb-impl</artifactId>
   <version>2.3.0.1</version>
  </dependency>
  <dependency>
   <groupId>javax.activation</groupId>
   <artifactId>activation</artifactId>
   <version>1.1.1</version>
  </dependency>


  <!-- https://mvnrepository.com/artifact/org.glassfish.jersey.media/jersey-media-json-jackson -->
  <dependency>
   <groupId>org.glassfish.jersey.media</groupId>
   <artifactId>jersey-media-json-jackson</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-grizzly2-http</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-grizzly2-servlet</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-jdk-http</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-simple-http</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-jetty-http</artifactId>
   <version>2.27</version>
  </dependency>

  <dependency>
   <groupId>org.glassfish.jersey.containers</groupId>
   <artifactId>jersey-container-jetty-servlet</artifactId>
   <version>2.27</version>
  </dependency>

  <!-- https://mvnrepository.com/artifact/org.glassfish.jersey.inject/jersey-hk2 -->
  <dependency>
   <groupId>org.glassfish.jersey.inject</groupId>
   <artifactId>jersey-hk2</artifactId>
   <version>2.27</version>
  </dependency>

  <!-- https://mvnrepository.com/artifact/javax.ws.rs/javax.ws.rs-api -->
  <dependency>
   <groupId>javax.ws.rs</groupId>
   <artifactId>javax.ws.rs-api</artifactId>
   <version>2.1.1</version>
  </dependency>
 </dependencies>

 <build>
  <plugins>
   <plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-compiler-plugin</artifactId>
    <version>3.8.0</version>
    <configuration>
     <release>${java.version}</release>
    </configuration>
   </plugin>
  </plugins>
 </build>

</project>


We can use curl to post new students and to get the results back. To post new students, we can do as follows on the command line with curl. If you don't have Curl in the command line you may just skip this step:

curl -d "Joao" -X POST http://localhost:9998/students
curl -d "Joana" -X POST http://localhost:9998/students

to get the list back:

curl -X GET http://localhost:9998/students/

and the result is:

["Joao","Joana"]

Now, let's build a client that is able to do a post into and get from a service as well.

package is.project3.rest;

import java.util.List;

import javax.ws.rs.client.Client;
import javax.ws.rs.client.ClientBuilder;
import javax.ws.rs.client.Entity;
import javax.ws.rs.client.Invocation;
import javax.ws.rs.client.WebTarget;
import javax.ws.rs.core.MediaType;


public class MyRESTClient {
 public static void main(String[] args) {
        Client client = ClientBuilder.newClient();
        WebTarget webTarget = client.target("http://localhost:9998/students");

  webTarget.request().post(Entity.entity("Jose", MediaType.TEXT_PLAIN));

  Invocation.Builder invocationBuilder =  webTarget.request(MediaType.APPLICATION_JSON);
  @SuppressWarnings("unchecked")
  List<String> response = invocationBuilder.get(List.class);
  
  response.forEach(System.out::println);
 }
}

And we get our final result:

Joao
Joana
Jose

That's it!

Playing with Kafka Streams

According to its own site, "Kafka Streams is a client library for building applications and microservices, where the input and output data are stored in Kafka clusters". In simplified terms, Kafka is a publish-subscribe system oriented to streams processing. We can think of Kafka as a kind of crossroads, where information travels from one application to another. For example, from a large-scale distributed microservice application to a monitoring system.

In this post I assume that you have been able to start Kafka and that it is running on your localhost on port 9092, together with Zookeeper, which is available on port 2181. There are sites dedicated to running Kafka, so I will overlook that issue.

I will solve a couple of exercises:

  • Counting the occurrences of each key
  • Converting the output from Long to String
  • Reduce()
  • Materialized views
  • Windowed streams

Overview

In the end we want to have the following arrangement:



Here, a producer is writing content to the topic, a dedicated application is reading the data from the topic (possibly in parallel with other applications) and outputting results to a second topic. In this second topic, we will have a dedicated consumer waiting to get the results computed by the streams application (in the middle).

To reach this setting, we start from the right. We can resort to shell applications provided with Kafka, to run a consumer waiting on the resultstopic topic. I did it with the following command:

kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic resultstopic

Note that the exact way of running this command depends on how you ran Kafka. It might happen that this fails to work, as I saw in some of the students' computers. In that case you might want to create a subscriber yourself. Please resort to other tutorials on Kafka to do that.

Before we reach the streams application, we will discuss the producer, because we need to see what the producer is going to put on the kstreamstopic. Let's reuse one that comes from the Kafka tutorials. Just don't run it yet, it will be the last piece of the puzzle:

package is.kafkastreamsclass;

//import util.properties packages
import java.util.Properties;

//import simple producer packages
import org.apache.kafka.clients.producer.Producer;

//import KafkaProducer packages
import org.apache.kafka.clients.producer.KafkaProducer;

//import ProducerRecord packages
import org.apache.kafka.clients.producer.ProducerRecord;

//Create java class named “SimpleProducer”
public class SimpleProducer {

 public static void main(String[] args) throws Exception{

  //Assign topicName to string variable
  String topicName = args[0].toString();

  // create instance for properties to access producer configs   
  Properties props = new Properties();

  //Assign localhost id
  props.put("bootstrap.servers", "localhost:9092");

  //Set acknowledgements for producer requests.      
  props.put("acks", "all");

  //If the request fails, the producer can automatically retry,
  props.put("retries", 0);

  //Specify buffer size in config
  props.put("batch.size", 16384);

  //Reduce the no of requests less than 0   
  props.put("linger.ms", 1);

  //The buffer.memory controls the total amount of memory available to the producer for buffering.   
  props.put("buffer.memory", 33554432);

  props.put("key.serializer", 
    "org.apache.kafka.common.serialization.StringSerializer");

  props.put("value.serializer", 
    "org.apache.kafka.common.serialization.LongSerializer");

  Producer<String, Long> producer = new KafkaProducer<>(props);

  for(int i = 0; i < 1000; i++)
   producer.send(new ProducerRecord<String, Long>(topicName, Integer.toString(i), (long) i));
  
  System.out.println("Message sent successfully to topic " + topicName);
  producer.close();
 }
}

This producer is somewhat dull, as it just sends a 1000 (key, value) pairs to the topic, where the key is a string and the value a long, but for now it will be enough for our experiments. I represent the output of the producer as follows:

"0" => 0
"1" => 1
"2" => 2
...
"999" => 999

Counting the occurrences of each key

Our focus here is the stream reader. Let us start by a basic one:

package is.kafkastreamsblog;

import java.io.IOException;
import java.util.Properties;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;


public class SimpleStreamsExercises {

 public static void main(String[] args) throws InterruptedException, IOException {
  String topicName = args[0].toString();
  String outtopicname = "resultstopic";

  java.util.Properties props = new Properties();
  props.put(StreamsConfig.APPLICATION_ID_CONFIG, "exercises-application");
  props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
  props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
  props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass());
    
  StreamsBuilder builder = new StreamsBuilder();
  KStream<String, Long> lines = builder.stream(topicName);

  KTable<String, Long> outlines = lines.
    groupByKey().count();
  outlines.toStream().to(outtopicname);
   
  KafkaStreams streams = new KafkaStreams(builder.build(), props);
  streams.start();
  
  System.out.println("Reading stream from topic " + topicName);
  
 }
}

Don't forget the pom.xml file:

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <groupId>is</groupId>
  <artifactId>kafkastreamsblog</artifactId>
  <version>0.0.1-SNAPSHOT</version>
  <packaging>jar</packaging>

  <name>kafkastreamsclass</name>
  <url>http://maven.apache.org</url>

 <properties>
  <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
  <java.version>10</java.version>
  <maven.compiler.source>${java.version}</maven.compiler.source>
  <maven.compiler.target>${java.version}</maven.compiler.target>
 </properties>

 <dependencies>
  <dependency>
   <groupId>junit</groupId>
   <artifactId>junit</artifactId>
   <version>3.8.1</version>
   <scope>test</scope>
  </dependency>
  <!-- https://mvnrepository.com/artifact/org.apache.kafka/kafka-clients -->
  <dependency>
   <groupId>org.apache.kafka</groupId>
   <artifactId>kafka-clients</artifactId>
   <version>2.0.0</version>
  </dependency>

  <!-- https://mvnrepository.com/artifact/com.fasterxml.jackson.core/jackson-databind -->
  <dependency>
   <groupId>com.fasterxml.jackson.core</groupId>
   <artifactId>jackson-databind</artifactId>
   <version>2.9.5</version>
  </dependency>

  <!-- https://mvnrepository.com/artifact/org.apache.kafka/kafka-streams -->
  <dependency>
   <groupId>org.apache.kafka</groupId>
   <artifactId>kafka-streams</artifactId>
   <version>2.0.0</version>
  </dependency>

 </dependencies>
 <build>
  <plugins>
   <plugin>
    <groupId>org.apache.maven.plugins</groupId>
    <artifactId>maven-compiler-plugin</artifactId>
    <version>3.8.0</version>
    <configuration>
     <release>10</release>
     <!-- <compilerArgs> <arg>add-modules</arg> <arg>javax.xml.bind</arg> 
      </compilerArgs> -->
    </configuration>
   </plugin>
  </plugins>
 </build>
</project>

To understand this application, we need to take a look at a diagram available on the Kafka Streams site, which summarizes the relation between the library's main classes:



The groupByKey operation converts the stream to a KGroupedStream, by creating records of values indexed by the keys. In our case, since the producer will output 1000 different keys, each key will have a single record (a 0 for key 0, a 1 for key 1, a value 2 for key 2 and so on). Hence, the following count() will compute 1 for all keys and convert the KGroupedStream into a KTable, which is similar to a regular database table, having the key as the primary key. This table will have 1000 registers, each with the value 1. To see the result, we cover the KTable back to a KStream using the toStream().

To run the experiment, we need to specify the topic where the streams application will receive the data, as a command line argument. The same for the producer. A lack to do this will crash the programs. In my case, I used kstreamstopic, but you may use another topic. Just start the applications in this order:

1 - Kafka-console-consumer.sh

2 - then, SimpleStreamsExercises

3 - finally, the Producer.

Regarding the issue of the order at which applications start, Kafka keeps messages on the topic for a configurable amount of time, so we could always get the messages from the topic later, if necessary. You may also repeat the execution of the Producer as many times as you want, with only slight changes in the results (the value of the count() will keep increasing).

Converting from Long to String

But if you run this setting you may end up getting nothing on the Kafka-console-consumer, except 1000 empty lines. Why? Because we are outputting longs instead of strings and, therefore, you will be looking at ASCII character 1.  Let us change our code slightly, to ensure that we can properly see the results of our operation:

StreamsBuilder builder = new StreamsBuilder();
KStream<String, Long> lines = builder.stream(topicName);

KTable<String, Long> outlines = lines.groupByKey().count();

outlines.mapValues(v -> "" + v).toStream().to(outtopicname, Produced.with(Serdes.String(), Serdes.String()));
   
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

What is new here? The mapValues(), which uses a lambda expression to transform the long value "v" into a String. However, this change alone crashes the program, because we specified the DEFAULT_VALUE_SERDES to be a Long. Hence, the attempt to write a String on the outtopicname Stream will crash the program. Therefore, we need te explicitly tell the library that we are producing the output with a String fromat (the Produced.with in the end). In other words, the final stream has a format that is different from the initial stream and from the KTable. These two had Long values, while the final stream has a String value.

Now, you should see this in the Kafka-console-consumer shell:

3
3
3
3
3
...

or whatever number of times you ran the whole application, instead of 3 (e.g., I'm actually seeing a thousand 10s).

This is still not very handy, because we cannot see the keys. To see them we may change the lambda expression in the mapValues to become:


mapValues((k, v) -> k + " => " + v)

i.e., it receives the key-value pair and replaces the value (because the function is "mapValues") by the string with the key, the arrow and the value, which is much nicer (don't worry about the 12 your case should be different, perhaps smaller):

...
989 => 12
990 => 12
991 => 12
992 => 12
993 => 12
994 => 12
995 => 12
996 => 12
997 => 12
998 => 12
999 => 12

Reduce()

What about summing all the values of a given key?

KTable<String, Long> outlines = lines.
    groupByKey().
    reduce((oldval, newval) -> oldval + newval);
outlines.mapValues((k, v) -> k + " => " + v).toStream().to(outtopicname, Produced.with(Serdes.String(), Serdes.String()));

Swap the count() by a reduce(). The lambda expression in the reduce keeps accumulating the new values that show up for the key. The reduce stores the result as it is stateful (mind the legend in the figure before, regarding the reduce()). For example, you might see the following output in Kafka-console-consumer shell:

985 => 3940
986 => 3944
987 => 3948
988 => 3952
989 => 3956
990 => 3960
991 => 3964
992 => 3968
993 => 3972
994 => 3976
995 => 3980
996 => 3984
997 => 3988
998 => 3992
999 => 3996

Materialized views

Now, the case for a materialized view. Materialized views actually allow us to query the tables, either directly by reaching for the value of a key, or in ranges, as show in this case:


package is.kafkastreamsblog;

import java.io.IOException;
import java.util.Properties;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.state.KeyValueIterator;
import org.apache.kafka.streams.state.QueryableStoreTypes;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;


public class SimpleStreamsExercises {

 private static final String tablename = "exercises";

 public static void main(String[] args) throws InterruptedException, IOException {
  String topicName = args[0].toString();
  String outtopicname = "resultstopic";

  java.util.Properties props = new Properties();
  props.put(StreamsConfig.APPLICATION_ID_CONFIG, "exercises-application");
  props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
  props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
  props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass());
    
  StreamsBuilder builder = new StreamsBuilder();
  KStream<String, Long> lines = builder.stream(topicName);

  KTable<String, Long> countlines = lines.
    groupByKey().
    reduce((oldval, newval) -> oldval + newval, Materialized.as(tablename));
  countlines.mapValues(v -> "" + v).toStream().to(outtopicname, Produced.with(Serdes.String(), Serdes.String()));


  KafkaStreams streams = new KafkaStreams(builder.build(), props);
  streams.start();
  
  
  System.out.println("Press enter when ready...");
  System.in.read();
  while (true) {
   ReadOnlyKeyValueStore<String, Long> keyValueStore = streams.store(tablename, QueryableStoreTypes.keyValueStore());
   System.out.println("count for 355:" + keyValueStore.get("355"));
   System.out.println();
   // Get the values for a range of keys available in this application instance
   KeyValueIterator<String, Long> range = keyValueStore.range("880", "980");
   while (range.hasNext()) {
     KeyValue<String, Long> next = range.next();
     System.out.println("count for " + next.key + ": " + next.value);
   }
   range.close();
   Thread.sleep(30000);
  }  
 }
}

After sending a few more messages with the Producer, and pressing Enter, we get this result on the streams application:

Press enter when ready...

count for 355:355

count for 880: 880
count for 881: 881
count for 882: 882
count for 883: 883
count for 884: 884
count for 885: 885
count for 886: 886
count for 887: 887
count for 888: 888
count for 889: 889
count for 89: 89
count for 890: 890
count for 891: 891
count for 892: 892
count for 893: 893
count for 894: 894
count for 895: 895
count for 896: 896
count for 897: 897
count for 898: 898
count for 899: 899
count for 9: 9
count for 90: 90
count for 900: 900
count for 901: 901
count for 902: 902
count for 903: 903
count for 904: 904
count for 905: 905
count for 906: 906
count for 907: 907
count for 908: 908
count for 909: 909
count for 91: 91
count for 910: 910
count for 911: 911
count for 912: 912
count for 913: 913
count for 914: 914
count for 915: 915
count for 916: 916
count for 917: 917
count for 918: 918
count for 919: 919
count for 92: 92
count for 920: 920
count for 921: 921
count for 922: 922
count for 923: 923
count for 924: 924
count for 925: 925
count for 926: 926
count for 927: 927
count for 928: 928
count for 929: 929
count for 93: 93
count for 930: 930
count for 931: 931
count for 932: 932
count for 933: 933
count for 934: 934
count for 935: 935
count for 936: 936
count for 937: 937
count for 938: 938
count for 939: 939
count for 94: 94
count for 940: 940
count for 941: 941
count for 942: 942
count for 943: 943
count for 944: 944
count for 945: 945
count for 946: 946
count for 947: 947
count for 948: 948
count for 949: 949
count for 95: 95
count for 950: 950
count for 951: 951
count for 952: 952
count for 953: 953
count for 954: 954
count for 955: 955
count for 956: 956
count for 957: 957
count for 958: 958
count for 959: 959
count for 96: 96
count for 960: 960
count for 961: 961
count for 962: 962
count for 963: 963
count for 964: 964
count for 965: 965
count for 966: 966
count for 967: 967
count for 968: 968
count for 969: 969
count for 97: 97
count for 970: 970
count for 971: 971
count for 972: 972
count for 973: 973
count for 974: 974
count for 975: 975
count for 976: 976
count for 977: 977
count for 978: 978
count for 979: 979
count for 98: 98
count for 980: 980

This seems awkward, because the 98 shows up between the 979 and the 980, but keep in mind that the keys are strings.

Windowed streams

What if we want to restrict the results to the last x minutes, being x variable? In this case we should do as follows:
package is.kafkastreamsblog;

import java.io.IOException;
import java.util.Properties;
import java.util.concurrent.TimeUnit;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.kstream.TimeWindows;
import org.apache.kafka.streams.kstream.Windowed;


public class SimpleStreamsExercises {

 public static void main(String[] args) throws InterruptedException, IOException {
  String topicName = args[0].toString();
  String outtopicname = "resultstopic";

  java.util.Properties props = new Properties();
  props.put(StreamsConfig.APPLICATION_ID_CONFIG, "exercises-application");
  props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
  props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
  props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass());
    
  StreamsBuilder builder = new StreamsBuilder();
  KStream<String, Long> lines = builder.stream(topicName);

  KTable<Windowed<String>, Long> addvalues = lines.
    groupByKey().
    windowedBy(TimeWindows.of(TimeUnit.MINUTES.toMillis(1))).
    reduce((aggval, newval) -> aggval + newval, Materialized.as("lixo"));
  addvalues.toStream((wk, v) -> wk.key()).map((k, v) -> new KeyValue<>(k, "" + k + "-->" + v)).to(outtopicname, Produced.with(Serdes.String(), Serdes.String()));

  KafkaStreams streams = new KafkaStreams(builder.build(), props);
  streams.start();
  
  System.out.println("Reading stream from topic " + topicName);
  
 }
}


We are basically applying a window of 1 minute to the results, and therefore we may get:

988 => 988
989 => 989
990 => 990
991 => 991
992 => 992
993 => 993
994 => 994
995 => 995
996 => 996
997 => 997
998 => 998
999 => 999

I.e., the sum of the values of the last minute. You may play with this value and change it for 10 minutes for example. You will notice that the results might differ (even without the need to send new messages with the producer):

988 => 2964
989 => 2967
990 => 2970
991 => 2973
992 => 2976
993 => 2979
994 => 2982
995 => 2985
996 => 2988
997 => 2991
998 => 2994
999 => 2997

In fact several variants of windows exist, but I will not cover them here.

Saturday, October 6, 2018

Creating a MySQL Datasource in WildFly 14

To create a MySQL Datasource you first need te configure the MySQL driver, which doesn't come configured by default. To do this, I did the following. Remember to have the MySQL server running. At first don't start WildFly.

Firstly, follow this link to install the MySQL Driver in WildFly 14.

Secondly, you need to start WildFly 14 and access the localhost at port 9990: http://localhost:9990/.

Note that this will not work if you don't have a management user. If this URL gives you a WildFly text page with no management options, you most likely need to run the add-user.sh or add-user.bat in the bin directory of the WildFly installation, before proceeding.

Once you do that, you need to select the following option:



Then, a few dialogs will pop up, but in the end you need to have the following data. Please keep in mind that at some point you need to insert the username and the password you use to access the MySQL. You should also note the name of the database: proj2 in this case. For each different database you will need a different datasource. In the end you will have the opportunity to test the connection to the database. If it fails, check the output of WildFly.

Please mind the java:/MySqlDS JNDI name. This is the name that you must give in the persistence.xml file in the Java Persistence API project deployed in the same WildFly:

<jta-data-source>java:/MySqlDS</jta-data-source>