I am a software developer. I am passionate about technology in general and distributed systems, Performance, Security in particular. I work mainly with open source software, specialising in Java and Unix-based operating systems.
Wednesday, April 18, 2018
Thursday, March 8, 2018
Cassandra Indexing
How Cassandra works?
Cassandra all about indexing. Good for write heavy.
Cassandra can be thought of as a key-value database it actually contains a lookup key for every data in the form of a primary key
Primary Key: (mandatory for create/read)
- Primary Key is mandatory
- Primary key is immutable.
PRIMARY KEY (partition_key):
A partition key will always belong to one node in a cluster and that partition’s data will always be found on that node.


We can never access the second-level data (for instance, the email of a user) without accessing the primary username key first.
CQL (Cassandra Query Language) to get email would be
SELECT "email" FROM "user_tweets" WHERE "username" = 'jochasinga'; //Primary key mandatory to pass along with query
Different other form of Primary Key:
PRIMARY KEY ((partition_key), primary_key1, primary_key2)
Primary keys 1 & 2 are also called cluster columns
Where the partition key is important for data locality (points to a node in a cluster), the clustering column specifies the order that the data is arranged inside the partition
Association looks like this in a map PARTITION KEY -> PRIMARY KEY -> DATA
Map [String, Map[String, Data]]. //keys of a map are unique
More information
SELECT * FROM "user_tweets" WHERE "email" = 'jo.chasinga@gmail.com'; // Throws Error. Reason no primary key passed along with query.
Cassandra introduced secondary index in order to address above query get email without primary key.
Query using secondary index (with no primary key) will span across all nodes in a cluster. (FanOut)
Cassandra supports secondary indexes on column families where the column name is known (Not on dynamic columns)
More Information on Secondary Index:https://pantheon.io/blog/cassandra-scale-problem-secondary-indexes
Important links:
Advantages with Cassandra:
- Supports write operations at massive scale
- Highly Available
- Elastic scalability (new nodes can be added to cluster seamlessly)
Key information on Cassandra
- Schema has to define as per Cassandra documentation
- Choosing right primary key and clustered keys will be a challenge for MO data set.
- Values of Primary keys in cassandra has to be immutable. ( update not allowed on primary key )
- Mandatory to provide partition key and cluster key for read operation
- Secondary indexes wont scale with Cassandra as they work with out partition key (Fan out pattern)
- Secondary indexes applicable only on known columns (static columns)
- Aggregations in Cassandra are not supported by the Cassandra nodes
Thursday, January 25, 2018
Thursday, December 21, 2017
Kafka Performance Tuning
We all know the power and advantages of KAFKA. Apache Kafka is publish-subscribe messaging system which basically has three major components
KAFKA CONSUMER
KAFKA PRODUCER and
KAFKA BROKER
Broker Side Configuration
Producer side
Compression
Batch size
Sync or Async
Consumer Side configuration
Fetch Size
Applicable for all consumer instance in a consumer group.
Broker Side Configuration
num.replica.fetchers
This configuration parameter defines the number of threads which will be replicating data from leader to the follower. Value of this parameter can be modified as per availability of thread. If we have threads available we should have more number of replica fetchers to complete replication in parallel.
replica.fetch.max.bytes
This parameter is all about how much data you want to fetch from any partition in each fetch request. It’s good to increase value for this parameter so that it helps to create replica fast in the followers
replica.socket.receive.buffer.bytes
In case of less thread available for creating replica, we can increase the size of buffer. It will help to hold more data if replication thread is slow as compared to the incoming message rate.
num.partitions
This is the very important configuration which we should be taken care while having Kafka in live. As many partitions are there, we can have that level of parallelism and write data in parallel which will automatically increase the throughput
Having more partitions slows down performance and throughput if the system OS configuration can’t capable of handle it.
Creating more partitions for a topic depends on available threads and disk.
num.io.threads
Setting value for I/O threads directly depends on how much disk you have in your cluster. These threads are used by server for executing request. We should have at least as many threads as we have disks.
Compression
compression.codec
Compression reduces disk footprint leading to faster reads and writes.Currently Kafka supports - Values are 'none', 'gzip' and 'snappy'
Property values are none, gzip and snappy.
Batch Size
Batch.size measures batch size in total bytes instead of the number of messages. It controls how many bytes of data to collect before sending messages to the Kafka broker. Set this as high as possible, without exceeding available memory. The default value is 16384.
If you increase the size of your buffer, it might never get full. The Producer sends the information eventually, based on other triggers, such as linger time in milliseconds.
Batch size always confusing what batch size will be optimal. Large batch size may be great to have high throughput but you might feel latency issue in that. So, we can conclude that latency and throughput is inversely proportional to each other.
Async Producers
Publish the message and get the callback to get the acknowledgement of send data status.
'producer.type=1' to make producer async
'queue.buffer.max.ms =duration of match window
'batch.num.messages" = number of messages to be sent in batch.
Large Messages
Consider placing large files on the standard storage and using Kafka to send a message with the file location. In many cases this can be much faster than using Kafka
to send the large file itself.
Sunday, March 19, 2017
Monitoring & Alerting for Micro Services
I got opportunity to design and develop Monitoring and Alerting framework for all the micro services deployed in the organisations. Monitoring majorly classified into 3
- System Monitoring
- Application Monitoring
- Server Monitoring
Any abnormolity on any of the above 3 will raise an alert. Alert can be "Slack Notification", "Email", "Pager" and "Call".
NOTE: This is not for micro services tracing. For micro services tracing sophisticated open source tools available like Jaeger and zipkin
NOTE: This is not for micro services tracing. For micro services tracing sophisticated open source tools available like Jaeger and zipkin
Monitoring & Alerting Tools used:
Technologies:
Fluentd - Used for data/log forwarding
OpenTSDB - Used for Time Series Data. Replaces old RRD tools
( competitors would be InfluxDB, Druid, Cassandra)
( competitors would be InfluxDB, Druid, Cassandra)
ElasticSearch - Used for log search (Ex: "ERROR" count > 1 on app1 raise an alert)
Monitoring Types:
System Monitoring:
includes CPU, Disk, I/O, Processes,Virtual Memory, DHCP, Network etc.
Application Monitoring :
includes Failed Services, Batch/Cron jobs, Cache monitoring,
DB monitoring, transaction, 3rd party interactions etc.
Server Monitoring:
Apache Tomcat, Ngnix, HAProxy , Request/Response latency, Server health etc
Scripts has to write for all of the above monitoring modules. Scripts can be written in Ruby where Sensu has many sensu plugins which will help to have less number of lines in scripts.
Deployment Topology:
All scripts has to upload to chef server with version. Project configuration and other artefacts will be uploaded to chef-server.
Chef-client can be run from development machine/laptop.

Micro service will be up and running with all monitoring scripts, configuration files, application jar,fluentd, opentsdb agents.
Below shows final Java micro services which will be up and running.

Detail explanation:
Flow - (a)
Sensu agent runs on Micro service. Agent periodically executes scripts(& server monitoring) and sends output to sensu server.
Sensu Server aggregates data from the micro services. Sensu forwards alerts to PagerDuty if any threshold breach by a sensu metric.
Uchiwa is the dashboard for Sensu. It gives nice alerts view in order Data centre, VMs, Metrics.
Use Cases: CPU, Disk, File I/O, Server health etc.
Flow - (b)
Fluentd agent runs on micro service is a data forwarder. Fluentd listens on a application log file path and forwards data to Elastic Search.
ElasticSearch does indexing for the log data. Error count cron job runs on elastic search which does search on "error" count on logs group by application. If there is any error on the log script forwards to sensu server which in turn converts to Pager Duty alert.
Use Cases: Server access logs( tells how many 401 Rest codes group by Region and application, Application logs error count group by Region and application)
Flow - (c)
This will be very interesting use case. I have used lot of time series metrics forwarded to OpenTSDB. But I would like to mention metric which helped a lot. REST call requests are recorded as time series metrics. Example: (Rest EndPoint, timestamp, hitCount) .
Use Cases:
For every rest call OpenTSDB client on micro-service sends data to OpenTSDB server. OpenTSDB graph shows traffic group by Region and Date. Which helps to understand how HTTP traffic on each data centre.
Example:
How to use sensu checks.
Below few checks on RAM & Disk
sensu_check 'ram' do
command 'check-ram.rb -w 20 -c 10'
interval node['monitor']['check_interval']
subscribers %w[all]
end
sensu_check 'disk' do
command 'check-disk.rb'
interval node['monitor']['check_interval']
subscribers %w[all]
end
Above ruby script raises alert if Ram breaches threshold. File "check-ram.rb" will be available at Checf cookbooks as default file
https://github.com/sensu-plugins/sensu-plugins-memory-checks/blob/master/bin/check-ram.rb
Subscribe: is for notification.
Sunday, October 23, 2016
Data pipeline design for Mobile Data Traffic using AWS
I worked for a Mobile Operating System Company which has 50 Million + user base across the globe. All these mobile generates a billion data per day.
What kind of Mobile data ?
Every user activity on the mobile is treated as a data. Below are few examples
a) Install/Uninstall app
b) Open an app
c) Closing an app
d) Total time spend on an app
e) Network connectivity details
f) Heart beat (Contains OS build no, single sim/dual sim, network operator names)
Problem:
Organisation needs data collection, data pipeline and analytics for all the mobile data traffic.
Solution:
I have leveraged AWS services(EC2, Route 53, ELB, EBS, S3, Lambda, Redshift) for implementation. Not just pipe line design and analytics also implemented robust monitoring and alerting system for the entire pipeline. Also took approaches which will minimise operational cost of AWS. Lot of factors tested(Apache Benchmark testing) and optimised while implementation.
Overall Architecture
STAGE - 1 (DATA Collection & Data Enriching)
Data Collection (Shopvac): Code name is ShopvacShopvac service should be front facing, low latency and high throughput. Shopvac Service is hosted on a EC2. Below is the network topology diagram before i go deep about service.
Mobiles:
All mobiles post data using HTTP POST on a host where DNS resolves host to a AWS Route 53.
AWS Route 53:
Route 53 does simple redirection to AWS ELB.
AWS ELB:
ELB has list of services where java process is running for data collection. In Non Amazon world the same has to be done using Zookeeper/Etcd for Service discovery.
Shopvac Service: (Data Collection)
Its a light weight java process running. It has a REST END Point which will be listening for mobile POST data. End point stores data to local file system ( /var/app/shopvac/metric/<eventname.json>) .
This file system acts as a buffer. FluentD will listen on this path and forward data to Amazon S3.
FluentD forwards when ever file size reaches 100MB size or 5 secs.
Fluentd is an Open Source log forwarder like Logstash, SysD, CollectD, RSyslog etc.
Fluentd parses data for event name and creates a file with event name if not exist otherwise append data to the file.
Below Shopvac Service insight
Every service developed in organisation has to be bootstrapped with chef with above process running.
Java Process(Vertx) : Used Vertx Async programming. Vertex has Rest End point which will get invoked on mobile POST. Data will be enriched for example stores City, State and Country information(gets from Latitude and Longitude) to the existing metrics.
Java process stores data on local file system.
FluentD: Fluentd listens on local files system where metrics are stored. Fluentd forwards mobile data to AWS S3. Also forwards logs to ElasticSearch.
OpenTSDB Client: All time series data are written to OpenTSDB. One most important data we store in OpenTSDB is REST End point call. When ever REST point is invoked it is stored in time series data. This will give insight how many times Rest end pint is invoked per day, per week, per hour this will give traffic insight. What duration has peak traffic on the server.
Sensu Client : Sensu scripts will be deployed along with chef. Basic scripts like Disk utilisation, RAM utilisation, Server health a lot other metrics sense will report to sense server. On Abnormality or threshold breach Sense Server will alert through Pager Duty.
STAGE - 2 (Data Processing)
AWS-S3 will be source of truth. All types of metrics will be forwarded blindly to S3.AWS Lambda (Serverless Architecture)
will listen on S3 bucket and filters for interested metrics. Forwards interested metrics to RedShift. Instead of directly forwarding to Redshift forwards to CMET service.(developed as a proxy service internally).
Note: Lambda will be charged based on cpu cycles it spent with the code. So code on lambda should be as less as possible and also should be error free. It is difficult to debug Lambda.
Note(Non-Amazon World):
S3 has to replace with Kafka
Lambda replace with Kafka Consumers/Storm
STAGE-3 (DATA Storing)
Selected Redshift as Storage because of its advantages.Lambda can directly forward metrics to Redshift but this will not be a right approach. As Data Warehouse will be exposed to events which are abnormal in behaviour and also every DB comes with concurrent DB connections at a time. There is a need for intermittent service which acts as a proxy to Redhisft. So CMET is a service which acts as a proxy and do connection management and dropping long holding connections.
- Long running queries will be stopped and will be marked as failure. So that lambda will retry. This helps when Redshift can't respond nor process at that point of time. Also helps reduce load on Redshift.
- Exposing read connections(GUI Visualization) for set of users and write connections for users like AWS like Lambda and Admin.
Saturday, January 31, 2015
HTTPS
Below content contains
1) What is HTTPs
2) Server Authentication
3) Important check list to follow for having secure website with good performance.
(web site links to check your server security online. Moazilla standards for configuring TLS on server)
4) Reasons for insecure communications over https.
Even static data should be encrypted. That's the best way to keep website secure.
In reality all the below 3 give secure website.
1)Authentication
2)Data Integrity (Data doesn't change between client and server)
3) Encryption (can any one see my conversation)
All the above 3 are taken care by Transport Layer Security(TLS)
HTTPS ==> HTTP running on top of TLS
HTTP (http running on top of TLS)
TLS
TCP
IP
Do we need to encrypt of all of the web data ?
My answer would be yes.
why we need to ?
For example casual surfing at restaurants with out https can give information to
hackers that what sites being visited if it is financ.yahoo.com. what shares
you are interested etc.
Hackers can change text, password etc if client doesn't connects to right Server. That's the reason we need Server Authentication.
1) If client wants to connect to the right server. Client has to connect over
https. Over https browser downloads server signed certificate(public key) this gives guarantee that client is connected
to right server.
Advantage of TLS
1) Passive and Active attackers cant listen in because we are encrypting the data.
2) Attacker cant tamper as data is check suming.
3) attackers cant impersonate.
Configuring TLS
1) Arent Certificate expensive
2) wont it make server and site slow ?
3) what are the configuration best practices
Important check list need to follow in the order
1) Get a 2048-bit TLS Certificate
2) Configure TLS on your servers.
3) Verify TLS server configuration
4) Monitor performance: resumption rate etc.
5) Tune Server configuration. Cache etc.
6) Investigate SPDY & HTTP2.0
1)Get a 2048-bit TLS Certificate
If there is any 1024-bit certificate on server better to migrate to 2048-bit.
Certificates are below types.
a) Free certificates ( which are for non commerfical use from StartSSL)
b) Single host ( google.com)
c) Multi-domain (google.co.in, google.co.us, google.co.uk)
d) Wildcard (*.mysite.com)
2) Configure TLS on your servers
More about Server Side TLS configuration in the blow link
https://wiki.mozilla.org/Security/Server_Side_TLS
3) Verify TLS server configuration
How to verify TLS Server configuration (Qualys provides online to test Server, browser etc)
https://www.ssllabs.com/ssltest/
It gives score and useful tips. Before you access any website you can check that site security aspects using this tool.
4) Monitor performance: resumption rate etc
Usually cryptography stuff consumes more CPU. Modern CPUs are designed to handle huge data traffic over TLS
Assymetric cryptography - verify the public certificate and do public crypto (This one is expensive)
Symetric Cryptography - how we encrypt the application data
5) Tune Server configuration
Using HTTP Keep alive and session resumption doesn't require full handshake. So handshake doesn't dominate CPU Usage.
.
6) Google developed a protocal SPDY which gives better page load performance over regular https connections.
SPDY1&2 not only improves client performance also does on Server. SPDY allows single connection to server instead of many
connections to server. Single connection means few handshakes, fewers sockets, few buffers to allocate. SPDY consumes less memory but
more CPU and also fewer worker threads.
Few more reasons for insecure communication.
Few reasons for broken cert between client and server. Developer points
1) Incorrect host name return by server in the cert.
2) Incomplete Certificate Chain
3) Expired Certificates.
Insecure references
Some secure websites having javascript/css code like below.
<script src="http://aaa.com/script.js"></script>
Some browsers wont allow(http:) type of communication. This script is blocked will not
execute. If browser allows also it is secruity leak.
Use Protocol relative URI's. Protocol relative uri's will be
<script src="//aaa.com/script.js"></script>
Even secure website can have insecure hrefs
<a href="http://abc.com"/>
use Protocol relative urls
<a href="//abc.com"/>
Insecure re directions are expensive
1) https -->redirect to --> http --> again redirect to https
HSTS (HTTP strict transport security) eliminates HTTP--> HTTPs redirects (costly operations)
Server can return with this header when returns a page.
Strict-Transport-Security: max-age=20491234; includeSubDomains
max-age in seconds. Remember this policy(HSTS) for this many seconds.
includeSubDomains is optional. says remember this policy for all the sub domains.
1) What is HTTPs
2) Server Authentication
3) Important check list to follow for having secure website with good performance.
(web site links to check your server security online. Moazilla standards for configuring TLS on server)
4) Reasons for insecure communications over https.
Even static data should be encrypted. That's the best way to keep website secure.
In reality all the below 3 give secure website.
1)Authentication
2)Data Integrity (Data doesn't change between client and server)
3) Encryption (can any one see my conversation)
All the above 3 are taken care by Transport Layer Security(TLS)
HTTPS ==> HTTP running on top of TLS
HTTP (http running on top of TLS)
TLS
TCP
IP
Do we need to encrypt of all of the web data ?
My answer would be yes.
why we need to ?
For example casual surfing at restaurants with out https can give information to
hackers that what sites being visited if it is financ.yahoo.com. what shares
you are interested etc.
Hackers can change text, password etc if client doesn't connects to right Server. That's the reason we need Server Authentication.
1) If client wants to connect to the right server. Client has to connect over
https. Over https browser downloads server signed certificate(public key) this gives guarantee that client is connected
to right server.
Advantage of TLS
1) Passive and Active attackers cant listen in because we are encrypting the data.
2) Attacker cant tamper as data is check suming.
3) attackers cant impersonate.
Configuring TLS
1) Arent Certificate expensive
2) wont it make server and site slow ?
3) what are the configuration best practices
Important check list need to follow in the order
1) Get a 2048-bit TLS Certificate
2) Configure TLS on your servers.
3) Verify TLS server configuration
4) Monitor performance: resumption rate etc.
5) Tune Server configuration. Cache etc.
6) Investigate SPDY & HTTP2.0
1)Get a 2048-bit TLS Certificate
If there is any 1024-bit certificate on server better to migrate to 2048-bit.
Certificates are below types.
a) Free certificates ( which are for non commerfical use from StartSSL)
b) Single host ( google.com)
c) Multi-domain (google.co.in, google.co.us, google.co.uk)
d) Wildcard (*.mysite.com)
2) Configure TLS on your servers
More about Server Side TLS configuration in the blow link
https://wiki.mozilla.org/Security/Server_Side_TLS
3) Verify TLS server configuration
How to verify TLS Server configuration (Qualys provides online to test Server, browser etc)
https://www.ssllabs.com/ssltest/
It gives score and useful tips. Before you access any website you can check that site security aspects using this tool.
4) Monitor performance: resumption rate etc
Usually cryptography stuff consumes more CPU. Modern CPUs are designed to handle huge data traffic over TLS
Assymetric cryptography - verify the public certificate and do public crypto (This one is expensive)
Symetric Cryptography - how we encrypt the application data
5) Tune Server configuration
Using HTTP Keep alive and session resumption doesn't require full handshake. So handshake doesn't dominate CPU Usage.
.
6) Google developed a protocal SPDY which gives better page load performance over regular https connections.
SPDY1&2 not only improves client performance also does on Server. SPDY allows single connection to server instead of many
connections to server. Single connection means few handshakes, fewers sockets, few buffers to allocate. SPDY consumes less memory but
more CPU and also fewer worker threads.
Few more reasons for insecure communication.
Few reasons for broken cert between client and server. Developer points
1) Incorrect host name return by server in the cert.
2) Incomplete Certificate Chain
3) Expired Certificates.
Insecure references
Some secure websites having javascript/css code like below.
<script src="http://aaa.com/script.js"></script>
Some browsers wont allow(http:) type of communication. This script is blocked will not
execute. If browser allows also it is secruity leak.
Use Protocol relative URI's. Protocol relative uri's will be
<script src="//aaa.com/script.js"></script>
Even secure website can have insecure hrefs
<a href="http://abc.com"/>
use Protocol relative urls
<a href="//abc.com"/>
Insecure re directions are expensive
1) https -->redirect to --> http --> again redirect to https
HSTS (HTTP strict transport security) eliminates HTTP--> HTTPs redirects (costly operations)
Server can return with this header when returns a page.
Strict-Transport-Security: max-age=20491234; includeSubDomains
max-age in seconds. Remember this policy(HSTS) for this many seconds.
includeSubDomains is optional. says remember this policy for all the sub domains.
Subscribe to:
Posts (Atom)






