Sunday, November 29, 2015

Data consistency in distributed systems: From ACID to BASE

Brewer's CAP conjecture proven by Gilbert and Lynch established that it is impossible to achieve consistency, high availability and partition tolerance together.  This has led to the design of distributed systems that provide weaker consistency guarantees while ensuring high availability even with node and communication failures - BASE (Basically Available Soft state Eventually consistent).  In traditional databases, the two-phase commit protocol guarantees consistency with updates either being committed to all of the N nodes or none at all.  Waiting for all N nodes to respond increases latency for the operation and also impacts availability,  since if any one of the N nodes fails to respond the update cannot occur.  In eventually consistent systems like Dynamodb, updates are considered complete once they are written to a subset of the nodes. They are eventually propagated to all the nodes, so it is possible that a subsequent read operation against a node may return stale data if the update has not yet propagated to that node.  As a result, it is possible for these systems to have different "versions" of the same data and hence they must also be able to reconcile these different versions. DynamoDb uses vector clocks to order the versions when possible - so latest wins - otherwise, relying on the client to reconcile conflicting versions.   It is possible to avoid conflicts entirely by using consensus algorithms like PAXOS or RAFT that ensure that only one entity can perform an update. The general consensus problem is about reaching agreement across a set of distributed processes that can fail due to faults in the infrastructure (network or node failure) or for malicious reasons (byzantine failures).  Most practical implementations of these algorithms adopt a leader/master based approach where all writes are always made to the master and then propagated to replicas as in MongoDb.  When the master fails a leader election process is initiated that elects a new master while demoting the failed master.  MongoDb supports various consistency vs latency tradeoffs by allowing you to specify different write-concerns: un-acknowledged, journaled (written to journal on master), or replica acknowledged (written to one or more replicas in addition to master).  Jepsen (a tool that tests partition tolerance of distributed systems) testing with MongoDb shows even with write-concern set to acknowledge writes to majority of the replicas in a cluster, you can still end up with missing writes! 

Sunday, August 10, 2014

Cluster management tools

Cluster management tools provide abstractions to run software applications on a collection of hardware (physical machines or VMs in a cloud). The tools allow you to declaratively specify the resources - e.g. tasks, services - required by the application. The mapping of resources to processes on the hardware is the responsibility of the tools.

Kubernetes

Kubernetes currently provides the following abstractions: pods, replication controllers and services. Pods are collections of containers that are akin to application virtual hosts. Services are used to setup proxies pointing to pods. The pods a service proxies are identified by labels.  Replication controllers are used to monitor a population of pods using labels. Kubernetes will ensure that the specified number of replicas are available. Kubernetes does not yet allow you to specify resource requirements (CPU, memory) for the application, but there are plans to support this.  This talk by Brandon Burns provides a good introduction to Kubernetes.  Kubernetes is still under active development.

Mesos

Mesos implements sharing of computing resources using resource offers. A resource offer is a list of free resources on multiple slaves. The decision about which resources to use are made by the programs sharing the resources offered by Mesos.  Mesos does not collect any resource requirements from these programs, instead the programs can “reject” offers. Mesos resources use OS isolation mechanisms like Linux containers which allow resources to be constrained by CPU and memory.  The isolation mechanism is pluggable via external isolator plugins like the Diemos plugin for Docker.  Diemos allows Mesos to use Docker containers with a specific image along with cpu and memory constraints as resources.  Mesosphere is a startup offering mesos on ec2 among other platforms. Currently, there does not seem to be any auto-scaling support for a Mesos cluster on ec2 - you preallocate ec2 instances to the cluster and Mesos will offer resources (e.g docker containers) from the set of ec2 instances.  This talk on Mesos for Cluster Management by Ben Hindman is a good introduction to Mesos.  For an academic perspective, see this Mesos paper from Berkeley.

There are several schedulers available for Mesos: Aurora, Marathon and soon .. Kubernetes-Mesos.

Aurora

Apache aurora was open-sourced by Twitter. Job definitions are written in python.
  import os 
  hello_world_process = Process(name = 'hello_world', cmdline = 'echo hello world')

 hello_world_task = Task(
  resources = Resources(cpu = 0.1, ram = 16 * MB, disk = 16 * MB),
  processes = [hello_world_process])

 hello_world_job = Job(
  cluster = 'cluster1',
  role = os.getenv('USER'),
  task = hello_world_task)

 jobs = [hello_world_job]

 Jobs are submitted using a cmd line client:
 aurora create cluster1/$USER/test/hello_world hello_world.aurora

Marathon

Marathon - by Mesosphere. It has a well documented REST API to create, start, stop and scale tasks. When creating tasks, you specify the number of instances to run and Marathon will ensure that those instances are available on the Mesos cluster.  Marathon allows for constraints to indicate that tasks should run only on certain slaves etc. You would have to run a load-balancer (e.g. HAProxy) on each host to proxy traffic from outside to the tasks.Jobs are defined in JSON format:

{
    "container": {
    "image": "docker:///libmesos/ubuntu",
    "options" : []
  },
  "id": "ubuntu",
  "instances": "1",
  "cpus": ".5",
  "mem": "512",
  "uris": [ ],
  "cmd": "while sleep 10; do date -u +%T; done"
}
Submit the job to Marathon via the REST API:
curl -X POST -H "Content-Type: application/json" localhost:8080/v2/apps -d@ubuntu.json

Kubernetes-Mesos

Kubernetes-Mesos - Work has just started for running Kubernetes pods on Mesos.

Omega

Google has been working on the problem of scheduling jobs on a cluster to maximize utilization.  Here is an overview of their work on the Omega scheduler.

Other tools

Apache Helix (https://github.com/linkedin/helix/) - opensourced by LinkedIn.
Autoscaling on EC2: http://aws.amazon.com/autoscaling/

Sunday, May 4, 2014

Meteor, load balancing and sticky sessions

Meteor clients establish a long-lived connection with the server that is uniquely identified by a session identifier to support the DDP protocol.  The DDP protocol allows Meteor clients to make RPC calls and also allows the server to keep the client updated with changes to data, i.e., Mongo documents. Meteor uses sockjs, which provides a cross-browser web-socket like API which falls back on long-polling when web sockets are not supported by the server.

Meteor with sockjs long-polling

Consider the following setup: Nginx is acting as a load balancer that is load-balancing 2 or more meteor servers.  The Nginx server configuration looks like this:
upstream meteor_server_lp {
   server localhost:3000;
   server localhost:3001;
}
server {
        listen       8084;
        server_name  localhost;

        location / {
            proxy_pass  http://meteor_server_lp;
        }

}
This configuration of Nginx does not support web-sockets, so the Meteor clients will use long polling.  Since long polling re-establishes connections every so often due to connection timeouts, such a configuration will require sticky sessions to ensure that client is directed to the same server that they have previously established a connection with. A sockjs connection that is directed to the wrong server by the Nginx load balancer will fail with a 404 Not Found.  The solution is to compile Nginx with the sticky module, and modify the configuration to be sticky like this:
upstream meteor_server_lp {
    sticky;
   server localhost:3000;
   server localhost:3001;
}
server {
        listen       8084;
        server_name  localhost;

        location / {
            proxy_pass  http://meteor_server_lp;
        }

}
When using this setup with AWS load balancers, sticky session needs to be enabled on the load balancer - app stickiness using the "route" cookie setup by Nginx's sticky module.

Meteor with web sockets

Nginx since 1.3.13, supports web-sockets using the protocol switch mechanism in HTTP/1.1. The configuration file looks like this:
upstream meteor_server {
   server localhost:3000;
   server localhost:3001;
}
server {
        listen       8082;
        server_name  localhost;

        location / {
            proxy_pass  http://meteor_server;
            proxy_http_version 1.1;
            proxy_set_header Upgrade $http_upgrade;
            proxy_set_header Connection "upgrade";
        }
}
This works as is without the need for any sticky sessions. The websocket connection between the client and any one of the two servers will be established when the client first connects or reconnects or if one of the servers go down.  AWS load balancers don't support websockets with http listeners, but it works with a tcp listener setup.  But this means any SSL termination must occur on the instances being load balanced. 

Friday, March 16, 2012

HAProxy and SSL

HAProxy does not have support for SSL. Common solution is to use Stud to handle SSL and send un-encrypted data to the backends.
Terminating SSL in the load balancer is not considered a good idea because it does not scale.
It is considered better to use webservers like Nginx with session caching enabled.
Good benchmark comparing Nginx, Stud and Stunnel is here- http://vincent.bernat.im/en/blog/2011-ssl-benchmark.html.
Another benchmark comparing stud,stunnel and nginx: http://matt.io/entry/uq and the follow up which establishes Nginx to be just as performant as Stud - the key is picking the right cipher.
http://matt.io/technobabble/hivemind_devops_alert:_nginx_does_not_suck_at_ssl/ur

Sunday, February 19, 2012

Tomcat with HAProxy/Nginx

Tomcat is usually fronted with a http server for various reasons - security, load balancing and additional functionality like URL-rewriting. Most common options for the proxy include: HTTPD, HAProxy and NGINx.

Compile HAProxy from source
$ make
$ make TARGET=generic
$ sudo make install

Resources:
http://www.tomcatexpert.com/blog/2010/07/12/trick-my-proxy-front-tomcat-haproxy-instead-apache
http://www.mulesoft.com/tomcat-proxy-configuration
http://haproxy.1wt.eu/download/1.2/doc/architecture.txt

Tuesday, January 31, 2012

Comet technology

Server Push, long polling, Good descriptions here: http://code.google.com/p/google-web-toolkit-incubator/wiki/ServerPushFAQ
Maturity of Comet implementations: http://cometdaily.com/maturity.html
Best Comet/Streaming server: Caplin Liberator (http://www.caplin.com/caplin_liberator.php)

Sunday, August 14, 2011

Building C++/.NET apps with MSBuild 4.0

In .NET 4.0/VS 2010, Microsoft replaced vcbuild.exe with msbuild.exe.
To build both .NET managed as well as native C++ apps, you only need .NET 4 along with Windows 7 SDK:
http://www.microsoft.com/download/en/details.aspx?displayLang=en&id=8279
There is no need to install VS 2010.
Here is a walkthrough for a simple hello world C++ app:
http://msdn.microsoft.com/en-us/library/dd293607.aspx
With VS 2010, vcbuild.exe is no longer used to build C++ projects.

For VS 2008 solution files, you will need Microsoft Windows 7 SDK and .NET 3.5:
http://www.microsoft.com/download/en/details.aspx?displaylang=en&id=3138
After installation, use the CMD shell (Programs->Windows 7 SDK->Cmd) to invoke msbuild on solution files.
This version of MSBuild (3.xx) uses vcbuild.exe to build C++ projects.

There is no need to install VS 2008.

Thursday, July 28, 2011

Event loop approach to concurrency

Event loop approach to concurrency as an alternative to threading - everything is non-blocking and executed via callbacks. The event loop is executing a queue of callbacks forever. This works well as long as the callbacks complete quickly! If a callback is going to take long it should fork another process. The primary application is networking - non-blocking I/O.
Douglas Crockford's presentation on event loop approach to concurrency

Libraries that use this approach include node.js, Ruby's Event Machine and Python's Twisted. and Java's new JDK7 Asynchronous IO and the older NIO library. The main difference between new Asychronous I/O and the older NIO - for NIO you are notified when the read operation is ready to start (data is available); while in Asychronous I/O - you are notified only when the read is completed (all data is read).

The design pattern being employed in all of these is the Reactor Design pattern. The Reactor pattern allows for the activation of handlers when events occur (e.g. activates handler to read data from socket when the data is available)

Wednesday, July 27, 2011

Five minute rule

The new five-minute rule
Compares the cost of holding data in memory vs disk I/O. With flash memory prices becoming cheaper, you can now pool memory from different machines to provide an ocean of RAM with low-latency.
RAM -> Flash memory -> Disk

Tuesday, July 26, 2011

CAP theorem

Eric Brewer's presentation of CAP theorem at PODC (Principles of Distributed Computing) keynote address: Consistency, Availability and Partition to network tolerance - only two of these properties can be possessed by shared data systems.
Consistency + Availability: Single-site databases (2-phase commit)
Consistency + Partitions: Distributed databases (pessimistic locking)
Availability + Partitions: DNS (conflict resolution)
Formally proven in 2002 paper by Seth Gilbert and Nancy Lynch.
BASE (Basically Available, Soft-state, Eventually consistent) is the opposite of ACID.

Verner Vogels article on Eventual Consistency
Great write on CAP here

Tuesday, July 19, 2011

Javascript and OOP

OOP is defined by three things: encapsulation, polymorphism and inheritance. Douglas Crockford's article claims it supports all three so it is an object-oriented language.

var AnimalClass = function Animal(name) {
this.name = name;

// private method
function sayPrivate() {
return "sayPrivate";
};

this.sayPrivileged = function() {
return sayPrivate();
}
}

// public method is added to the prototype
AnimalClass.prototype.say = function (something) {
return this.name + something;
}

var anAnimal = new AnimalClass("foo");
alert(anAnimal.name);
alert(anAnimal.say("ha"));
alert(anAnimal.sayPrivileged("ha"));

Typical way to implement inheritance in Javascript is via object-augmentation. For example, the underscore library defines the following function to extend any given object with the properties of the passed in object.
// Extend a given object with all the properties in passed-in object(s).
_.extend = function(obj) {
each(slice.call(arguments, 1), function(source) {
for (var prop in source) {
if (source[prop] !== void 0) obj[prop] = source[prop];
}
});
return obj;
};

_.extend(anAnimal, { "foo" : "bar" });
alert(anAnimal.foo);

Monday, July 18, 2011

Monitoring app performance

Front end performance: Speedtracer - can tell you how much time was spent on - DOM processing, garbage collection in the browser.
New relic tracks both front-end and backend performance by injecting javascript into the brower:
http://blog.newrelic.com/2011/05/17/how-rum-works/

Friday, July 15, 2011

Spring Roo 1.1.5 with GWT & GAE

mkdir rooapp
cd rooapp
start roo shell
roo> project --topLevelPackage com.xxx.rooapp --java 6
roo>persistence setup --provider DATANUCLEUS --database GOOGLE_APP_ENGINE
roo>entity --class ~.model.Product --testAutomatically
roo>field string --fieldName name --notNull
roo>field string --fieldName id --notNull
roo>field date --fieldName dateIntroduced --type java.util.Date --notNull
roo>field number --type java.lang.Float --fieldName unitPrice --notNull
roo>field string --fieldName description --notNull
roo>web gwt setup
I encountered a number of problems ...
The POM file it generated setup gae:home to be in the maven repo. Had to change it for things to work.
Compile failures with gwt:compile:
[INFO] [ERROR] Line 3: The import com.xxx.rooapp.server.gae.UserServiceLocator cannot be resolved
[INFO] [ERROR] Line 4: The import com.xxx.rooapp.server.gae.UserServiceWrapper cannot be resolved
This because GWT needs access to these sources - these files were not in the GWT source path!
Ran into lots of other errors around the generated classes for GAE and the GWT DesktopInjector ...
[ERROR] Generator 'com.google.gwt.inject.rebind.GinjectorGenerator' threw an exception while rebinding 'com.xxx.rooapp.client.scaffold.ioc.DesktopInjector'

Wednesday, July 13, 2011

Deploying apps with Puppet

Use Puppet master/agent

I want to deploy a simple application that installs a file in /tmp on a single box.

Here is the puppet module definition:
/etc/puppetlabs/puppet/modules/myapp/manifests/init.pp:
class myapp {
file { 'testfile' : path => '/tmp/testfile-local', ensure => 'present', content => 'Test', mode => 0640}
}

Typically, the puppet master uses a site.pp file for node definitions:
/etc/puppet/manifests/site.pp:

node development {
include "myapp"
}

node staging {
include "myapp"
}

Start the puppet master:
puppet master --no-daemonize --verbose
notice: Starting puppet master version 2.6.4

This is telling puppet master that myapp needs to be installed on the development box. If I have a puppet agent running on the client (which in this case happens to be the same box), then it will apply the latest configuration from the server automatically when it polls the next time.

I could run the agent manually onetime by ssh'ing into the box and invoking the agent:
[root@learn ~]# puppet agent --no-daemonize --onetime --server puppet --verbose
info: Retrieving plugin
info: Caching catalog for puppet
info: Applying configuration version '1310559260'
notice: /Stage[main]/Myapp/File[testfile-local]/ensure: created
notice: Finished catalog run in 0.02 seconds

CONS:
The main problem here is that I am not sure how you to achieve this via puppet master itself .. have not checked the UI - but surely there must be a way to update specific nodes - i.e. dev nodes only ?

You can always automate this using a capistrano script:
set :user, "root"
task :deploy
role :app, ‘development’
run 'puppet agent --no-daemonize --onetime --server puppet --verbose'
end

The problem with using a script like this is that the environment information has to be maintained in two places – in the script and the site.pp file. One way to avoid this is to generate this script from the site.pp file.

PROS:
You are using Puppet to install the application, just like a sysadmin would to ensure that machines were setup correctly. Easy to sell to ops.

Server-less puppet

Use capistrano to run puppet “apply” on the relevant nodes. Here is a Capfile for the app. It assumes that you have manifests checked out on the clients.

set :user, "root"

task :development do
role :app, “development”
end

task :staging do
role :app, "staging"
end

task :deploy do
# TODO: checkout manifests to module path
run "puppet apply -e \"include myapp\""
end

To deploy to development environment, you would run cap for that environment:

rg6977:puppet Thoughtworks$ cap development deploy
* executing `development'
* executing `deploy'
* executing "puppet apply -e \"include myapp\""
servers: ["192.168.56.101"]
Password:
[192.168.56.101] executing command
command finished

PROS:
Node definitions are now in Capistrano and not in puppet. Puppet is used only to install the application in a given environment. Puppet does a poor job of managing environment specific information. See http://docs.puppetlabs.com/guides/environment.html.

CONS:

Node definitions are now in Capistrano and not in puppet. Puppet is used only to install the application in a given environment.
This is more code than the previous approach. Puppet “apply” will only apply manifests to the local machine. So, if your application must be installed on a web server and db server, your Capistrano script needs to invoke puppet apply on the appropriate manifest (db vs web).

Tuesday, July 12, 2011

Amazon ec2

Command line api:
ec2-describe-images -o amazon
ec2-run-instances -k
ec2-describe-instances
ec2-authorize default -p 22 -open port for ssh
- ssh into the box
ssh -i root@xxx.amazonaws.com
ec2-terminate-instances
Good tutorial here

Puppet with Nagios

Puppet : Resources, aggregate resources using "defines" and "classes", organize using "modules" (see Puppet Language Guide)
Run-Stages used to control the order of resource management. The "sigils" - magical operators - <| and @ when doubled up - are very powerful.
Good article on configuring nagios with puppet. Also see the original Puppet example it references.

Thursday, May 26, 2011

SNA Projects Blog at LinkedIn

Building a terabyte-scale data cycle at LinkedIn with Hadoop and Project Voldemort:
http://project-voldemort.com/blog/2009/06/building-a-1-tb-data-cycle-at-linkedin-with-hadoop-and-project-voldemort/
There are many other interesting articles at http://sna-projects.com/blog/

Amazon Dynamo
http://www.allthingsdistributed.com/2007/10/amazons_dynamo.html

Running JBehave tests with xvfb

- Install xvfb
- Start xvfb: Xvfb :1 -screen 0 1024x768x24
- export DISPLAY=:1
- firefox (should start without errors)
- install vnc server and use a VNC client like ChickenOfVnc to connect and see your tests running.

Friday, January 28, 2011

Rails Architecture/Performance issues

Typical Rails architecture:
Nginx/Apache --> HAProxy --> Mongrel cluster

HAProxy is a http proxy that forwards requests from web server to an available Mongrel instance. queue requests up if all Mongrels are busy.

HTTP options: Passenger, Mongrel, Thin, Unicorn

Monitoring tools:
New Relic, top/iostat, God, Monit

Great blog on reasons for memory bloat in Rails:
http://www.engineyard.com/blog/2009/thats-not-a-memory-leak-its-bloat/
Github architecture: https://github.com/blog/530-how-we-made-github-fast
Rack:Bug: https://github.com/brynary/rack-bug/

Thursday, January 27, 2011

Mailing list managers (Listserv, Mailman) vs MTAs (Sendmail, Postfix)

Some terminology:
Mail Transfer Agent: MTA implements both the client (sending) and server (receiving) portions of the Simple Mail Transfer Protocol.
An MTA receives a message from another MTA, MSA or MUA. If the recipient mailbox is not hosted locally, then it routes it to another MTA. The Domain Name System (DNS) associates a mail server to a domain with mail exchanger (MX) resource records containing the domain name of a host providing MTA services.
MX: A mail exchanger record (MX record) is a type of resource record in the Domain Name System that specifies a mail server responsible for accepting email messages on behalf of a recipient's domain and a preference value used to prioritize mail delivery.
MUA: Mail user agent - this is an email client like GMail, Outlook etc. MUA uses POP3 (Post Office Protocol) or IMAP (Internet Message Access Protocol) to retrieve messages from an MTA.
MSA: Mail submission agent - sits between the MUA and MTA. Functionally same as MTA.

Sendmail/Postfix use information from the Domain Name System (DNS) to figure out which IP addresses go with which mailboxes.
1. Setup a domain name -e.g. companyA.com
2. Configure name servers for your domain (primary and secondary)
3. Configure MX records for your domain.
4. After the name servers are setup, register your domain using one of the registries.
5. Configure sendmail to listen for mail/route outgoing mail.
6. Mailing lists are configured by setting up aliases. In Postfix, edit the /etc/aliases file. It has the format:
alias: address1,address2. After updating this file, you are usually required to run commands to update the internal db file used by Postfix/Sendmail.

Mailing list managers like Mailman integrate with an MTA like Sendmail or Postfix, so that when new lists are created, or lists are removed, Postfix's alias database will be automatically updated.