• About Blog

    What's Blog?

    A blog is a discussion or informational website published on the World Wide Web consisting of discrete, often informal diary-style text entries or posts.

  • About Cauvery Calling

    Cauvery Calling. Action Now!

    Cauvery Calling is a first of its kind campaign, setting the standard for how India’s rivers – the country’s lifelines – can be revitalized.

  • About Quinbay Publications

    Quinbay Publication

    We follow our passion for digital innovation. Our high performing team comprising of talented and committed engineers are building the future of business tech.

Showing posts with label Scalable Applications. Show all posts
Showing posts with label Scalable Applications. Show all posts

Monday, July 25, 2022

Understand Code Refactoring Techniques to Improve your Code Quality

Code Refactoring Techniques
Photo by Danial Igdery

Now a days, agile teams are under tremendous pressure to write code faster with enhanced functionality in short time. There would be some or the other functionality added at the last moment or just before the release. As engineers are under pressure, the functionality gets implemented in a sloppy manner, which may technically work but may lead to dirty code.

Dirty code usually results from a developer’s inexperience, shortcuts taken to meet increasingly tight deadlines, poor management of the project, several different developers working on a project over time, or some combination of all of the above.

Bad code comes at a price, and writing good code isn’t that complicated. Let's understand what's code refactoring.

Code Refactoring

In computer programming and software design, code refactoring is the process of restructuring existing computer code—changing the factoring—without changing its external behavior. Refactoring is intended to improve the design, structure, and/or implementation of the software (its non-functional attributes), while preserving its functionality. 

Potential advantages of refactoring may include improved code readability and reduced complexity; these can improve the source code's maintainability and create a simpler, cleaner, or more expressive internal architecture or object model to improve extensibility. 

Another potential goal for refactoring is improved performance; software engineers face an ongoing challenge to write programs that perform faster or use less memory.

Code refactoring helps to change this dirty code into clean code and it helps to make it easier to extend the code and add new features easily in the future. Also, helps to improve the more objective attributes of code such as code length, code duplication, and coupling and cohesion, all of which correlate with ease of code maintenance and your code will use less memory and perform better and faster.

How to Perform Code Refactoring?

Now that you know the answer to the question of what is code refactoring, and you know many of its potential benefits, how exactly do you refactor code?

There are many approaches and techniques to refactor the code. Let’s discuss some popular ones.

Red-Green Refactor

Red-Green is the most popular and widely used code refactoring technique in the Agile software development process. This technique follows the test-first approach to design and implementation, this lays the foundation for all forms of refactoring.

Red - The first step starts with writing the failing red-test. You stop and check what needs to be developed.

Green - In the second step, you write the simplest enough code and get the development pass green testing.

Refactor - Find ways to improve the code and implement those improvements, without adding new functionality.

Refactoring by Abstraction

This technique is mostly used by developers when there is a need to do a large amount of refactoring. Mainly we use this technique to reduce the redundancy (duplication) in our code. This involves class inheritances, hierarchy, creating new classes and interfaces, extraction, replacing inheritance with the delegation, and vice versa.

Pull-Up Method - It pulls code parts into a superclass and helps in the elimination of code duplication.

Push-Down Method - It takes the code part from a superclass and moves it down into the subclasses.

Refactoring by abstraction allows you to make big changes to large chunks of code gradually. In this way, you can still release the system regularly, even with the change still in progress.

Composing Method

Code that is too long is difficult to understand and difficult to implement. The composing method is a code refactoring approach that helps to streamline code and remove any code duplications. This is done through extraction and inline techniques.

Extraction - We break the code into smaller chunks to find and extract fragmentation. After that, we create separate methods for these chunks, and then it is replaced with a call to this new method. Extraction involves class, interface, and local variables.

Inline - Refactoring also helps to create simpler, more streamlined code. It helps to remove unnecessary methods within the code and replaces them with the content of the method. After that, we delete the method from our program.

Simplifying Methods

As legacy code gets older and older, it tends to become more polluted and complex. In this sense, simplifying methods help to simplify the logic. These methods include adjusting the interaction between different classes, along with adding a new parameter or removing and replacing certain parameters with explicit methods.

Simplifying Conditional Expressions - Conditional statement in programming becomes more logical and complicated over time. You need to simplify the logic in your code to understand the whole program. There are so many ways to refactor the code and simplify the logic. Some of them are: consolidate conditional expression and duplicate conditional fragments, decompose conditional, replace conditional with polymorphism, remove control flag, replace nested conditional with guard clauses, etc.

Simplifying Method Calls - In this approach, we make method calls simpler and easier to understand. We work on the interaction between classes, and we simplify the interfaces for them. Examples are: adding, removing, and introducing new parameters, replacing the parameter with the explicit method and method call, parameterize method, making a separate query from modifier, preserve the whole object, remove setting method, etc.

Extract Method

The extract method is one of the techniques for code refactoring that helps to decrease complexity, while also increasing the overall readability of the code.

When you find that a class has so many responsibilities and too much thing is going on or when you find that a class is unnecessary and doing nothing in an application, you can move the code from this class to another class and remove it completely from the existing class. It involves the moving of a fragment or block of code from its existing method and into a newly created method, which is clearly named in order to explain its function.

Conclusion

Engineers are the only ones responsible for writing good and quality code. We should all make it a habit to write good code from the very beginning. Writing clean code isn’t complicated and doing so will help both you and your colleagues. A clean and well-organized code is always easy to change, easy to understand, and easy to maintain.

Clean Code: A Handbook of Agile Software Craftsmanship by Robert C. Martin

Even bad code can function. Every year, countless hours and significant resources are lost because of poorly written code but it doesn't have to be that way.

Refactoring: Improving the Design of Existing Code by Martin Fowler

Any fool can write code that a computer can understand. Good programmers write code that humans can understand.


Quick Tips

  • Always perform code refactoring in small chunks by making your code slightly better and leaves the application in a working state. Run jUnit tests after making small changes in the refactoring process. Without running these tests, you create a risk of introducing new bugs.

  • Do not create any new features or functionality during the refactoring process. You should refactor the code before adding any updates or new features into your existing code.

  • Refactoring process always results in complete regression, don't forget to involve your QA team in the process.

Wednesday, October 20, 2021

Understand and Analyze Java Thread Dump

Thread Dump Image
Photo by Mel Poole

Microservices

Also known as the microservices architecture, is an architectural style that structures an application as a collection of services that are:
  • Easily Maintainable and Testable
  • Loosely Coupled
  • Independently Deployable
  • Organized around Business Capabilities
  • Owned by a Small Team

The microservices architecture enables the rapid, frequent and reliable delivery of large, complex applications. It also enables an organization to evolve its technology stack.

The decentralization of business logic increases the flexibility and most importantly decouples the dependencies between two or more components, this being one of the major reasons as to why many companies are moving from monolithic architecture to a microservices architecture.


What is a Thread?

All of us have probably written a program that displays "Hello World!!" or given word is a palindrome or not etc. These are sequential programs that have a beginning, an execution sequence and an end, at any given point of time during the execution of a program, there is a single point of execution.

A single thread is also similar, as it has a beginning, an execution sequence and an end. However, a thread itself is not a program, a thread cannot run on its own, it runs within a program.

A program can consist of many lightweight processes called threads. The real excitement surrounding threads is not about a single sequence. It helps to achieve parallelism wherein, a program is divided into multiple threads and results in better performance. All threads within a process share the same memory space and might have dependency on each other in some cases.


Lifecycle of a Thread

For understanding a thread dump in detail, it is essential to know all the states a thread passes through during its lifecycle. A thread can assume one of these following states at any given point of time:

NEW

Initial state of a thread when we create an instance of Thread or Runnable. It remains in this state until the program starts the thread.

RUNNABLE

The thread becomes runnable after a new thread is started. A thread in this state is considered to be executing its task.

BLOCKED

A thread is in the blocked state when it tries to access an object that is currently locked by some other thread. When the locked object is unlocked and hence available for the thread, the thread moves back to the runnable state.

WAITING

A thread transitions to the waiting state while waiting for another thread to perform a task and transitions back to the runnable state only when another thread signals the waiting thread to resume execution.

TIMED_WAITING

A timed waiting state is a thread waiting for a specified interval of time and transitioning back to the runnable state when that time interval expires. The thread is waiting for another thread to do some work for up to a specified waiting time.

TERMINATED

A runnable thread enters the terminated state after it finishes its task.


Thread Dumps

A thread dump contains a snapshot of all the threads active at a particular point during the execution of a program. It contains all relevant information about the thread and its current state.

A new age application development involves multiple numbers of threads. Each thread requires certain resources, performs certain task related to the program. This can boost the performance of an application as threads can utilize available CPU cores. But we do have some trade-offs, for example, sometimes multiple threads may not co-ordinate well with each other and a deadlock situation may arise depending on the program. So, if something goes wrong, we can use thread dumps to identify the state of our threads.

As Java has been most popular language among application development, let's consider our application is built using spring-boot. If you want to take a snapshot of application threads, then we can go ahead with taking thread dump. A JVM thread dump is a listing of the state of all threads that are part of the process at that particular point of time. It contains information about the thread’s stack with other important information. The dump will be in a plain text format, the contents can be saved and analysis can be done either manually or using some UI that are available.

Analysis of thread dumps can help in following areas:
  • Tweak JVM performance
  • Tweak application performance
  • Identify threads related problems within application.

Now we know the basics of thread and it's life cycle. Let's get into the next stage where we will explore how to take a thread dump from any running Java application. 

There are multiple ways to take a thread dumps. Am going to discuss about some JVM based tools and can be executed from the CLI or GUI tools.

Java Stack Trace

One of the easy way to generate a thread dump is by using jStack. jStack is a utility that ships with JVM, it can be used from the CLI and it expects the PID of the process for which we want to generate the thread dump.

jstack -l 1129 > thread_dump.txt


Java Command

JCMD is a command-line utility that ships with the JDK and are used to send diagnostic command requests to the JVM, where these requests are useful for controlling Java Flight Recordings, troubleshoot, and diagnose JVM and Java Applications. It must be used on the same machine where the JVM is running, and have the same effective user and group identifiers that were used to launch the JVM.

We can use the Thread.print command of jcmd to get a list of thread dumps for a particular process specified by the PID.

jcmd 1129 Thread.print > thread_dump.txt


Java Console

The jconsole GUI is a monitoring tool that complies to the Java Management Extensions (JMX) specification. It ships with the JDK and uses the extensive instrumentation of the Java VM to provide information about the performance and resource consumption of applications running on the Java platform.

Using the jconsole tool, we can inspect each thread’s stack trace when we connect it to a running java process. Then, in the Thread tab, we can see the name of all running threads. To detect a deadlock, we can click on the Detect Deadlock in the bottom right of the window. If a deadlock is detected, it will appear in a new tab otherwise a No Deadlock Detected will be displayed.

To launch the GUI tool, just type the below command on CLI.

jconsole


VisualVM

VisualVM is a GUI tool that helps us troubleshoot, monitor and profile Java applications. It perfectly fits all requirements of application developers, system administrators, quality engineers and end users.

As it's an external program, you need to download and install it on your machine. The GUI tool is very easy to use and lot of things you can monitor and troubleshoot things related to Java Applications.


Understanding Thread Dump Contents

Now, Let’s see what are the things we can explore using thread dumps. If we observe the thread dump, we can see a lot of information. However, if we take one step at a time, it can be fairly simple to understand.

1129:
2021-10-13 12:57:15
Full thread dump Java HotSpot(TM) 64-Bit Server VM (25.261-b12 mixed mode):

"Attach Listener" #142 daemon prio=9 os_prio=31 tid=0x00007f8dc7146000 nid=0x440b waiting on condition [0x0000000000000000]
   java.lang.Thread.State: RUNNABLE

"http-nio-8080-Acceptor" #138 daemon prio=5 os_prio=31 tid=0x00007f8dc7fab800 nid=0x9c03 runnable [0x000070000c59f000]
   java.lang.Thread.State: RUNNABLE
        at sun.nio.ch.ServerSocketChannelImpl.accept0(Native Method

  • The thread dump entry shown above, starts with the name of the thread Attach Listener whose ID is 142 thread (indicated by#142) created by the JVM after the application has started.
  • The daemon keyword after the thread number indicates that it's a daemon thread, which means that it will not prevent the JVM from shutting down if it is the last running thread.
  • After that we have are less important pieces of metadata about the thread like a priority, os priority, thread identifier, and native identifier.
  • The last piece of information is the most important, the state of the thread and its address in the JVM. The thread can be in one of the states as explained earlier in thread life cycle.

I am sure that most of us may not want to analyze the thread dump in plain text file. One can use GUI tools to analyse thread dumps. 


Conclusion

Now you know, what's thread dump and how it can be generated. Also it's useful in understanding and diagnosing problems in multithreaded applications. With proper knowledge, regarding the thread dumps and it's structure, the information contained in dump etc,  can be utilized to identify the root cause of the problems quickly.


Monday, September 13, 2021

Redis Overview and Benchmark

What is Redis?
Image Courtesy Morioh

ReDiS which stands for Remote Directory Server, is an open source in-memory data store, used as a database and as a cache. Redis provides data structures such as strings, hashes, lists, sets and sorted sets. Redis has built-in replication, Lua scripting, LRU eviction, transactions, and different levels of on-disk persistence and provides high availability via Redis Sentinel and automatic partitioning with Redis Cluster.

Redis is an open source, advanced key-value store and an apt solution for building highperformance, scalable web applications.

Redis has three main features that sets it apart from others:

  • Redis holds its database entirely in the memory, using the disk only for persistence.
  • Redis has a relatively rich set of data types when compared to many key-value data stores.
  • Redis can replicate data to any number of slaves.


Following are certain advantages of Redis:

  • Exceptionally fast − Redis is very fast and can perform about 110000 SETs per second, about 81000 GETs per second.
  • Supports rich data types − Redis natively supports most of the datatypes that developers already know such as list, set, sorted set, and hashes. This makes it easy to solve a variety of problems as we know which problem can be handled better by which data type.
  • Operations are atomic − All Redis operations are atomic, which ensures that if two clients concurrently access, Redis server will receive the updated value.
  • Multi-utility tool − Redis is a multi-utility tool and can be used in a number of use cases such as caching, messaging-queues (Redis natively supports Publish/Subscribe), any short-lived data in your application, such as web application sessions, web page hit counts, etc.

Redis Monitoring

Availability, the redis server will respond to the PING command when it's running smoothly.

$ redis-cli -h 127.0.0.1 ping
PONG

Cache Hit Rate

This information can be calculated with the help of INFO command.

$ redis-cli -h 127.0.0.1 info stats | grep keyspace
keyspace_hits:1069963628
keyspace_misses:2243422165

Workload Statistics

The first two stats talks about connections and commands processed where last two stats talk about bytes received and sent from the redis server.

$ redis-cli -h 127.0.0.1 info stats | grep "^total"
total_connections_received:1687889
total_commands_processed:5602955422
total_net_input_bytes:198210899161
total_net_output_bytes:309040592973

Key Space

Anytime to know number of keys in the database, use this command. The size of the keyspace with a quick drop or spike in the number of keys is a good indicator of issues.

$ redis-cli -h 127.0.0.1 info keyspace
# Keyspace
db0:keys=3857884,expires=277,avg_ttl=259237

Clear Keys

We can clear all the keys from the Redis, using the below command.

$ redis-cli -h 127.0.0.1
127.0.0.1:6379> flushall


How to Perform Redis Benchmark?

Redis benchmark is the utility to check the performance of Redis by running n commands simultaneously.

redis-benchmark [option] [option value]

Option Description
-h Specifies server host name 127.0.0.1
-p Specifies server port 6379
-c Specifies number of parallel connections, default is 50
-n Specifies total number of requests, default is 100000
-d Specifies data size of SET/GET value in bytes, default is 3
-r Use random keys for SET/GET/INCR
-q Forces Quiet to Redis. Just shows query/sec values
-l Generates loop, Run the tests forever
-t Only runs the comma-separated list of tests
--csv
Output in CSV format

$ redis-benchmark -h 127.0.0.1 -n 100000 -q
PING_INLINE: 57306.59 requests per second
PING_BULK: 57273.77 requests per second
SET: 56657.22 requests per second
GET: 57012.54 requests per second
INCR: 57240.98 requests per second
LPUSH: 57045.07 requests per second
RPUSH: 56657.22 requests per second
LPOP: 57142.86 requests per second
RPOP: 57175.53 requests per second
SADD: 56369.79 requests per second
HSET: 55679.29 requests per second
SPOP: 54704.60 requests per second
LPUSH (needed to benchmark LRANGE): 52798.31 requests per second
LRANGE_100 (first 100 elements): 35448.42 requests per second
LRANGE_300 (first 300 elements): 17618.04 requests per second
LRANGE_500 (first 450 elements): 12812.30 requests per second
LRANGE_600 (first 600 elements): 10036.13 requests per second
MSET (10 keys): 47281.32 requests per second

References

Tuesday, June 15, 2021

Materialized Views in RDBMS - Is it a View or Table?

 Materialized View Image

Google uses structured data to understand the content on the page and use that data to display in richer features in the search results.

Below you can see the difference in the search result data provided by Google for the OnePlus 8T product from official site vs amazon site. The amazon one has additional data rendered in the same search results.

Google Search Results
Image Courtesy - Shushma

As we started to build the feature to feed the product details to Google, the data has to be retrieved from various tables from the catalog schema based on multiple  conditions.

Once we had the API ready, we did load testing, having 25 threads with a ramp-up time of 5 seconds and test duration of 30 seconds. Though we know the search engine bots won't put so much load on the site but the performance of the API was a concern to check.

Performance Test Result 01

If you look at the 95 percentile of time, it took 10 seconds to respond and the number of requests reached 430 only. Team members came up with the idea of caching the data on redis, so that next call will be faster. Caching at application layer for this use case will not be of that great help as we have a TTL of 600 seconds on Redis.

I recalled in one of my previous assignments, where we had a similar use case and we had utilised the caching of the data on the database (materialized view) side instead on the application layer which had given better results.

Understand Terminologies

Let's start with TABLE, it's basically an organized storage for your data in the form of rows and columns. You can easily query the table using predicates on the columns. 

To simplify the queries or maybe to apply different security mechanisms on data being accessed you can use VIEWs. Think of a view as glasses through which you can look at your data without knowing the actual tables and realtionship details etc.

If the table is a storage, a view is just a way of looking at it or a projection of the data as view doesn't store the data physically. If you query a table, you fetch its data directly. On the other hand, when you query a view, you are basically querying another query that is stored in the view's definition. But the query planner is aware of that and build a plan to merge the two together and give the results.

Between the table and view, we have the MATERIALIZED VIEW. It's a view that has a query in its definition and uses this query to fetch the data directly from the storage, but it also has it's own storage that basically acts as a cache in between the underlying table(s). 

Now you might have a question, if any data is updated in the underlying table(s) whether they will reflect in materialized view. The simple answer is NO as materialized view acts like a cache, it wont reflect the changes. We need to refresh the materialized view, through a process that would cause it's definition's query to be executed again on actual data and rebuilds the cache.

Materialized View Approach

We created a materialized view to cache the data on the database side, as we were retrieving the data from multiple tables. The required data is available in it, so we can avoid lot of computation at the real-time and also the fetched data won't change very frequently.

With the code changes, we did same load testing having 25 threads with a ramp-up time of 5 seconds and test duration of 30 seconds.

Performance Test Result 02

If you look at the 95 percentile of time it took 75 milli seconds to respond and the number of requests reached to 8376. The performance gain with respect to throughput is 19x and 95 percentile of time is 134x time faster. The results looks wonderful and it's without caching on the application layer side.

Query Plan

  • The original query which used to join multiple tables with various conditions had the query planning time of 0.852 ms and execution time of 5.511 ms.
  • The query which we fire on materialized view from the application layer had the query planning time of 0.087 ms and execution time of 0.052 ms.

Conclusion

  • The Materialized View is supported by all relational databases like Oracle, MySQL, Postgres etc.
  • The  Materialized View is a powerful tool enabling many performance improvements while providing another way of ensuring data consistency.

Thursday, January 21, 2021

Rate Limiter Implementation — Sliding Log Algorithm

Sliding Log Image

API Rate Limiting

Rate limiting is a strategy to limit the access to APIs. It restricts the number of API calls that a client can make within any given timeframe. This helps to defend the API against abuse, both unintentional and malicious scripts.

Rate limits are often applied to an API by tracking the IP address, API keys or access tokens, etc. As an API developers, we can choose to respond in several different ways when a client reaches the limit.

  • Queueing the request until the remaining time period has elapsed.
  • Allowing the request immediately but charging extra for this request.
  • Most common one is rejecting the request (HTTP 429 Too Many Requests)

Sliding Log Algorithm

Sliding Log rate limiting involves tracking a time stamped log for each consumer request. These logs are usually stored in a hash set or table that is sorted by time. Logs with timestamps beyond a threshold are discarded. When a new request comes in, we calculate the sum of logs to determine the request rate. If the request would exceed the threshold rate, then it is held.

The advantage of this algorithm is that it does not suffer from the boundary conditions of fixed windows. The rate limit will be enforced precisely and because the sliding log is tracked for each consumer, you don’t have the rush effect that challenges fixed windows. However, it can be very expensive to store an unlimited number of logs for every request. It’s also expensive to compute because each request requires calculating a summation over the consumers prior requests, potentially across a cluster of servers. As a result, it does not scale well to handle large bursts of traffic or denial of service attacks.

Please refer to the Understanding Rate Limiting Algorithms blog where the Sliding Log and other algorithms have been explained in detail.

Building a Springboot Application with API Rate Limiter

Create a new spring boot application from Spring Initializr with dependency on spring web module.

Unzip the downloaded project and import to your IDE. We are going to implement a simple calculator REST APIs that can do operations like add and subtract.

@RestController
@RequestMapping(value = "/api/calculator")
public class CalculatorController {
    @GetMapping(value = "/add")
    public ResponseEntity<Calculator> add(@RequestParam int left, @RequestParam int right) {
        return ResponseEntity.ok(Calculator.builder()
                .operation("add").answer(left + right).build());
    }
    @GetMapping(value = "/subtract")
    public ResponseEntity<Calculator> subtract(@RequestParam int left, @RequestParam int right) {
        return ResponseEntity.ok(Calculator.builder()
                .operation("subtract").answer(left - right).build());
    }
}

Let’s ensure that our above APIs are up and running as expected. You can use the cURL or PostMan to make an API call.

curl -X GET -H "Content-Type: application/json" 'http://localhost:9090/api/calculator/add?left=20&right=30'{"operation":"add","answer":50}

Now that we have APIs ready to consume, next let’s introduce some subscription plans with rate limits. Let’s assume that we have the following subscription plans for our clients:
Free Subscription allows 2 requests per 60 seconds.
Basic Subscription allows 10 requests per 60 seconds.
Professional Subscription allows 20 requests per 60 seconds.

Each API client gets a unique API key that they must send along with each request. This would help us identify the client and subscription plan linked.

public enum SubscriptionPlan {

    SUBSCRIPTION_FREE(2, 60),
    SUBSCRIPTION_BASIC(10, 60),
    SUBSCRIPTION_PROFESSIONAL(20, 60);

    private final int requestLimit;
    private final int windowTime;

    SubscriptionPlan(int requestLimit, int windowTime) {
        this.requestLimit = requestLimit;
        this.windowTime = windowTime;
    }

    public int getRequestLimit() {
        return this.requestLimit;
    }

    public int getWindowTime() {
        return this.windowTime;
    }

}

Next we create a subscription service which will store the references for each of the API client in a memory.

@Service
public class SubscriptionService {

    private final Map<String, UserRequestData>
            subscriptionCacheMap = new ConcurrentHashMap<>();

    public UserRequestData resolveSubscribedUserData(String subscriptionKey) {
        return subscriptionCacheMap.computeIfAbsent(
                subscriptionKey, this::resolveUser);
    }

    private UserRequestData resolveUser(String subscriptionKey) {
        if (subscriptionCacheMap.containsKey(subscriptionKey)) {
            return subscriptionCacheMap.get(subscriptionKey);
        }
        return buildUserLog(
                resolveSubscriptionPlanByKey(subscriptionKey));
    }

    private UserRequestData buildUserLog(SubscriptionPlan subscriptionPlan) {
        return new UserRequestData(
                subscriptionPlan.getRequestLimit(), 
                subscriptionPlan.getWindowTime());
    }

    private SubscriptionPlan resolveSubscriptionPlanByKey(String subscriptionKey) {
        if (subscriptionKey.startsWith("PS1129-")) {
            return SubscriptionPlan.SUBSCRIPTION_PROFESSIONAL;
        } else if (subscriptionKey.startsWith("BS1129-")) {
            return SubscriptionPlan.SUBSCRIPTION_BASIC;
        }

        return SubscriptionPlan.SUBSCRIPTION_FREE;
    }

}

Let’s understand the implementation. The API client sends an API key with the X-Subscription-Key request header. We use the SubscriptionService to get the user reference for the API key and check whether the request is allowed or not with the help of methods.

In order to enhance the client experience of the API, we will add the following additional response headers to send information about the rate limit.
  • X-Rate-Limit-Remaining — number of tokens remaining in the current time window.
  • X-Rate-Limit-Retry-After-Seconds — remaining time in seconds until the bucket is refilled with new tokens.
We can call UserRequestData methods getRequestWaitTime and getRemainingRequests, to get the count of the remaining requests and the time remaining until the next sliding log respectively. The implementation provided in this class is self explanatory and easy to understand the same.

public class UserRequestData {

    private int requestLimit;
    private int windowTimeInSec;
    private Queue<Long> requestTimeStamps;

    public UserRequestData(
            int requestLimit, int windowTimeInSec) {
        this.requestLimit = requestLimit;
        this.windowTimeInSec = windowTimeInSec;
        this.requestTimeStamps =
                new ConcurrentLinkedDeque<Long>();
    }

    public int getRemainingRequests() {
        return requestLimit - requestTimeStamps.size();
    }

    public int getRequestWaitTime() {
        long currentTimeStamp =
                System.currentTimeMillis() / 1000;
        int initialElapsedTime =
                (int) (currentTimeStamp - requestTimeStamps.peek());
        return initialElapsedTime > windowTimeInSec
                ? 0 : windowTimeInSec - initialElapsedTime;
    }

    public boolean isServiceCallAllowed() {
        long currentTimeStamp =
                System.currentTimeMillis() / 1000;
        evictOlderRequestTimeStamps(currentTimeStamp);

        if (requestTimeStamps.size() >= this.requestLimit) {
            return false;
        }

        requestTimeStamps.add(currentTimeStamp);
        return true;
    }

    public void evictOlderRequestTimeStamps(long currentTimeStamp) {
        while (requestTimeStamps.size() != 0 &&
                (currentTimeStamp - requestTimeStamps.peek() 
                        > windowTimeInSec)) {
            requestTimeStamps.remove();
        }
    }

}

Here is the implementation of the Interceptor to validate the request with rate limiter to see whether we accept or reject the request.

@Component
public class RateLimiterInterceptor implements HandlerInterceptor {

    private static final String
            HEADER_SUBSCRIPTION_KEY = "X-Subscription-Key";
    private static final String
            HEADER_LIMIT_REMAINING = "X-Rate-Limit-Remaining";
    private static final String
            HEADER_RETRY_AFTER = "X-Rate-Limit-Retry-After-Seconds";
    private static final String
            SUBSCRIPTION_QUOTA_EXHAUSTED =
            "You've exhausted your API Request Quota. " +
            "Please upgrade your subscription plan.";

    @Autowired
    private SubscriptionService subscriptionService;

    @Override
    public boolean preHandle(HttpServletRequest request,
                             HttpServletResponse response,
                             Object handler) throws Exception {
        String subscriptionKey =
                request.getHeader(HEADER_SUBSCRIPTION_KEY);
        if (StringUtils.isEmpty(subscriptionKey)) {
            response.sendError(HttpStatus.BAD_REQUEST.value(),
                    "Missing Request Header: " +
                       HEADER_SUBSCRIPTION_KEY);
            return false;
        }

        UserRequestData userRequestData = subscriptionService
                        .resolveSubscribedUserData(subscriptionKey);
        if (!userRequestData.isServiceCallAllowed()) {
            int waitTime = userRequestData.getRequestWaitTime();
            response.addHeader(HEADER_RETRY_AFTER,
                    String.valueOf(waitTime));

            response.setContentType(
                    MediaType.APPLICATION_JSON_VALUE);
            response.sendError(
                    HttpStatus.TOO_MANY_REQUESTS.value(),
                    SUBSCRIPTION_QUOTA_EXHAUSTED);
            return false;
        }

        response.addHeader(HEADER_LIMIT_REMAINING,
                String.valueOf(
                    userRequestData.getRemainingRequests()));
        return true;
    }
}

Finally, let’s add the interceptor to the InterceptorRegistry of Springboot so that the RateLimitInterceptor intercepts each request to our calculator API endpoints.

@SpringBootApplication
public class SlidingWindowApplication implements WebMvcConfigurer {

    @Autowired
    @Lazy
    private RateLimiterInterceptor interceptor;

    public void addInterceptors(InterceptorRegistry registry) {
        registry.addInterceptor(interceptor)
                .addPathPatterns("/api/calculator/**");
    }

    public static void main(String[] args) {
        SpringApplication.run(SlidingWindowApplication.class, args);
    }

}

Let invoke calculator API to see the behaviour.

curl -X GET 'http://localhost:9090/api/calculator/add?left=20&right=30'
{"timestamp":"2021-01-03T12:56:20.047+0000","status":400,"error":"Bad Request","message":"Missing Request Header: X-Subscription-Key","path":"/api/calculator/add"}

The client has to send the API key within the http header otherwise the interceptor will not process the request. Let’s add the API key to the header and make the call.

curl -v -X GET -H "X-subscription-key:A1129-12" 'http://localhost:9090/api/calculator/subtract?left=20&right=30'
* Connected to localhost (::1) port 9090 (#0)
> GET /api/calculator/subtract?left=20&right=30 HTTP/1.1
> Host: localhost:9090
> User-Agent: curl/7.64.1
> Accept: */*
> X-subscription-key:A1129-12
>
< HTTP/1.1 200
< X-Rate-Limit-Remaining: 1
< Content-Type: application/json
< Transfer-Encoding: chunked
< Date: Sun, 03 Jan 2021 12:57:09 GMT
<
* Connection #0 to host localhost left intact
{"operation":"subtract","answer":-10}
* Closing connection 0

You can see the API key is added in the header, the API responds to our request and also it has added response header which shows how many rate is remaining for the API key.

Let’s make 2 more calls then we should see that we exhausted our rate for the free plan and returns 429 as response.

curl -v -X GET -H "X-subscription-key:A1129-12" 'http://localhost:9090/api/calculator/subtract?left=20&right=30'
* Connected to localhost (::1) port 9090 (#0)
> GET /api/calculator/subtract?left=20&right=30 HTTP/1.1
> Host: localhost:9090
> User-Agent: curl/7.64.1
> Accept: */*
> X-subscription-key:A1129-12
>
< HTTP/1.1 429
< X-Rate-Limit-Retry-After-Seconds: 24
< Content-Type: application/json
< Transfer-Encoding: chunked
< Date: Sun, 03 Jan 2021 12:58:58 GMT
<
* Connection #0 to host localhost left intact
{"timestamp":"2021-01-03T12:58:58.176+0000","status":429,"error":"Too Many Requests","message":"You've exhausted your API Request Quota. Please upgrade your subscription plan.","path":"/api/calculator/subtract"}
* Closing connection 0

It looks like we have successfully implemented the rate limiter using the Sliding Log algorithm. We can keep adding endpoints and the interceptor would apply the rate limit for each request.

As usual, the source code for the above spring boot implementation is available over on GitHub.

Monday, January 4, 2021

Rate Limiter Implementation — Token Bucket Algorithm

Token Bucket Image

API Rate Limiting

Rate limiting is a strategy to limit the access to APIs. It restricts the number of API calls that a client can make within any given timeframe. This helps to defend the API against abuse, both unintentional and malicious scripts.

Rate limits are often applied to an API by tracking the IP address, API keys or access tokens, etc. As an API developers, we can choose to respond in several different ways when a client reaches the limit.

  • Queueing the request until the remaining time period has elapsed.
  • Allowing the request immediately but charging extra for this request.
  • Most common one is rejecting the request (HTTP 429 Too Many Requests)

Token Bucket Algorithm

Assume that we have a bucket, the capacity is defined as the number of tokens that it can hold. Whenever a consumer wants to access an API endpoint, it must get a token from the bucket. Token is removed from the bucket if it’s available and accept the request. If the token is not available then the server rejects the request.

As requests are consuming tokens, we also need to refill them at some fixed rate and time, such that we never exceed the capacity of the bucket. Let’s consider an API that has a rate limit of 100 requests per minute. We can create a bucket with a capacity of 100, and a refill rate of 100 tokens per minute.

Please refer to the Understanding Rate Limiting Algorithms blog where the Token Bucket and other algorithms have been explained in detail.

Building a Springboot Application with API Rate Limiter

Create a new spring boot application from Spring Initializr with dependency on spring web module.

Unzip the downloaded project and import to your IDE. Let’s begin by adding the bucket4j dependency to our pom.xml

<dependency>
    <groupId>com.github.vladimir-bukhtoyarov</groupId>
    <artifactId>bucket4j-core</artifactId>
    <version>4.10.0</version>
</dependency>

We are going to implement a simple calculator REST APIs that can do operations like add and subtract.

@RestController
@RequestMapping(value = "/api/calculator")
public class CalculatorController {

    @GetMapping(value = "/add")
    public ResponseEntity<Calculator> add(@RequestParam int left, @RequestParam int right) {
        return ResponseEntity.ok(Calculator.builder()
                .operation("add").answer(left + right).build());
    }

    @GetMapping(value = "/subtract")
    public ResponseEntity<Calculator> subtract(@RequestParam int left, @RequestParam int right) {
        return ResponseEntity.ok(Calculator.builder()
                .operation("subtract").answer(left - right).build());
    }

}

Let’s ensure that our above APIs are up and running as expected. You can use the cURL or PostMan to make an API call.

curl -X GET -H "Content-Type: application/json" 'http://localhost:9090/api/calculator/add?left=20&right=30'
{"operation":"add","answer":50}

Now that we have APIs ready to consume, next let’s introduce some subscription plans with rate limits. Let’s assume that we have the following subscription plans for our clients:
  • Free Subscription allows 2 requests per 60 seconds.
  • Basic Subscription allows 10 requests per 60 seconds.
  • Professional Subscription allows 20 requests per 60 seconds.

Each API client gets a unique API key that they must send along with each request. This would help us identify the client and subscription plan linked.

public enum SubscriptionPlan {

    SUBSCRIPTION_FREE(2),
    SUBSCRIPTION_BASIC(10),
    SUBSCRIPTION_PROFESSIONAL(20);

    private int bucketLimit;

    private SubscriptionPlan(int bucketLimit) {
        this.bucketLimit = bucketLimit;
    }

    public int getBucketLimit() {
        return this.bucketLimit;
    }

    public Bandwidth getBandwidth() {
        return Bandwidth.classic(bucketLimit,
                Refill.intervally(bucketLimit,
                        Duration.ofMinutes(1)));
    }

}

Next we create a subscription service which will store the bucket reference for each of the API client in a memory.

@Service
public class SubscriptionService {

    private final Map<String, Bucket>
            subscriptionCacheMap = new ConcurrentHashMap<>();

    public Bucket resolveBucket(String subscriptionKey) {
        return subscriptionCacheMap.computeIfAbsent(
                subscriptionKey, this::getSubscriptionBucket);
    }

    private Bucket getSubscriptionBucket(String subscriptionKey) {
        return buildBucket(
                resolveSubscriptionPlanByKey(subscriptionKey)
                        .getBandwidth());
    }

    private Bucket buildBucket(Bandwidth limit) {
        return Bucket4j.builder().addLimit(limit).build();
    }

    private SubscriptionPlan resolveSubscriptionPlanByKey(
            String subscriptionKey) {
        if (subscriptionKey.startsWith("PS1129-")) {
            return SubscriptionPlan.SUBSCRIPTION_PROFESSIONAL;
        } else if (subscriptionKey.startsWith("BS1129-")) {
            return SubscriptionPlan.SUBSCRIPTION_BASIC;
        }

        return SubscriptionPlan.SUBSCRIPTION_FREE;
    }
}

Let’s understand the implementation. The API client sends an API key with the X-Subscription-Key request header. We use the SubscriptionService to get the bucket for this API key and check whether the request is allowed by consuming a token from the bucket.

In order to enhance the client experience of the API, we will add the following additional response headers to send information about the rate limit.
  • X-Rate-Limit-Remaining - number of tokens remaining in the current time window.
  • X-Rate-Limit-Retry-After-Seconds - remaining time in seconds until the bucket is refilled with new tokens.
We can call ConsumptionProbe methods getRemainingTokens and getNanosToWaitForRefill, to get the count of the remaining tokens in the bucket and the time remaining until the next refill, respectively. The getNanosToWaitForRefill method returns 0 if we are able to consume the token successfully.

Let’s create a RateLimitInterceptor and implement the rate limit code in the preHandle method instead of writing in every API method as we will have cleaner implementation.

@Component
public class RateLimiterInterceptor implements HandlerInterceptor {

    private static final String
            HEADER_SUBSCRIPTION_KEY = "X-Subscription-Key";
    private static final String
            HEADER_LIMIT_REMAINING = "X-Rate-Limit-Remaining";
    private static final String
            HEADER_RETRY_AFTER = "X-Rate-Limit-Retry-After-Seconds";
    private static final String
            SUBSCRIPTION_QUOTA_EXHAUSTED =
            "You've exhausted your API Request Quota. " +
            "Please upgrade your subscription plan.";

    @Autowired
    private SubscriptionService subscriptionService;

    @Override
    public boolean preHandle(HttpServletRequest request,
                             HttpServletResponse response,
                             Object handler) throws Exception {
        String subscriptionKey =
                request.getHeader(HEADER_SUBSCRIPTION_KEY);
        if (StringUtils.isEmpty(subscriptionKey)) {
            response.sendError(HttpStatus.BAD_REQUEST.value(),
                    "Missing Request Header: " +
                        HEADER_SUBSCRIPTION_KEY);
            return false;
        }

        Bucket tokenBucket = subscriptionService
                .resolveBucket(subscriptionKey);
        ConsumptionProbe consumptionProbe =
                tokenBucket.tryConsumeAndReturnRemaining(1);
        if (!consumptionProbe.isConsumed()) {
            long waitTime =
                    consumptionProbe.getNanosToWaitForRefill()
                            / 1_000_000_000;
            response.addHeader(HEADER_RETRY_AFTER,
                    String.valueOf(waitTime));

            response.setContentType(
                    MediaType.APPLICATION_JSON_VALUE);
            response.sendError(
                    HttpStatus.TOO_MANY_REQUESTS.value(),
                    SUBSCRIPTION_QUOTA_EXHAUSTED);
            return false;
        }

        response.addHeader(HEADER_LIMIT_REMAINING,
                String.valueOf(
                    consumptionProbe.getRemainingTokens()));
        
        return true;
    }
}

Finally, let’s add the interceptor to the InterceptorRegistry of Springboot so that the RateLimitInterceptor intercepts each request to our calculator API endpoints.

@SpringBootApplication
public class TokenBucketApplication implements WebMvcConfigurer {

   @Autowired
   @Lazy
   private RateLimiterInterceptor interceptor;

   public void addInterceptors(InterceptorRegistry registry) {
      registry.addInterceptor(interceptor)
            .addPathPatterns("/api/calculator/**");
   }

   public static void main(String[] args) {
      SpringApplication.run(TokenBucketApplication.class, args);
   }

}

Let invoke calculator API to see the behaviour.

curl -X GET -H "Content-Type: application/json" 'http://localhost:9090/api/calculator/add?left=20&right=30'
{"timestamp":"2020-12-25T12:43:43.239+0000","status":400,"error":"Bad Request","message":"Missing Request Header: X-Subscription-Key","path":"/api/calculator/add"}

The client has to send the API key within the http header otherwise the interceptor will not process the request. Let’s add the API key to the header and make the call.

curl -v -X GET -H "Content-Type: application/json" -H "X-subscription-key:A1129-12" 'http://localhost:9090/api/calculator/add?left=20&right=30'
* Connected to localhost (::1) port 9090 (#0)
> GET /api/calculator/add?left=20&right=30 HTTP/1.1
> Host: localhost:9090
> User-Agent: curl/7.64.1
> Accept: */*
> Content-Type: application/json
> X-subscription-key:A1129-12
>
< HTTP/1.1 200
< X-Rate-Limit-Remaining: 1
< Content-Type: application/json
< Transfer-Encoding: chunked
< Date: Fri, 25 Dec 2020 12:46:06 GMT
<
* Connection #0 to host localhost left intact
{"operation":"add","answer":50}
* Closing connection 0

You can see the API key is added in the header, the API responds to our request and also it has added response header which shows how many rate is remaining for the API key.

Let’s make 2 more calls then we should see that we exhausted our rate for the free plan and returns 429 as response.

curl -v -X GET -H "Content-Type: application/json" -H "X-subscription-key:A1129-12" 'http://localhost:9090/api/calculator/add?left=20&right=30'
* Connected to localhost (::1) port 9090 (#0)
> GET /api/calculator/add?left=20&right=30 HTTP/1.1
> Host: localhost:9090
> User-Agent: curl/7.64.1
> Accept: */*
> Content-Type: application/json
> X-subscription-key:A1129-12
>
< HTTP/1.1 429
< X-Rate-Limit-Retry-After-Seconds: 51
< Content-Type: application/json
< Transfer-Encoding: chunked
< Date: Fri, 25 Dec 2020 12:49:11 GMT
<
* Connection #0 to host localhost left intact
{"timestamp":"2020-12-25T12:49:11.358+0000","status":429,"error":"Too Many Requests","message":"You've exhausted your API Request Quota. Please upgrade your subscription plan.","path":"/api/calculator/add"}
* Closing connection 0

It looks like we have successfully implemented the rate limiter using the Token Bucket algorithm. We can keep adding endpoints and the interceptor would apply the rate limit for each request.

As usual, the source code for the above spring boot implementation is available over on GitHub.

Featured Post

Your AI Sidekick: How Claude took over Pritee’s Repetitive tasks

  It was a classic Wednesday morning in our Bengaluru office . Pritee, one of our sharpest Project Managers, had just stepped out of a stake...