High-Level Design
In this section, we discuss query API design, data model, and high-level design.
Query API design
The purpose of the API design is to have an agreement between the client and the server. In a consumer app, a client is usually the end-user who uses the product. In our case, however, a client is the dashboard user (data scientist, product manager, advertiser, etc.) who runs queries against the aggregation service.
Let’s review the functional requirements so we can better design the APIs:
-
Aggregate the number of clicks of ad_id in the last M minutes.
-
Return the top N most clicked ad_ids in the last M minute.
-
Support aggregation filtering by different attributes.
We only need two APIs to support those three use cases because filtering (the last requirement) can be supported by adding query parameters to the requests.
API 1: Aggregate the number of clicks of ad_id in the last M minutes.
| API | Detail |
|---|---|
| GET /v1/ads/{:ad_id}/aggregated_count | Return aggregated event count for a given ad_id |
Table 1 API for aggregating the number of clicks
Request parameters are:
| Field | Description | Type |
|---|---|---|
| from | Start minute (default is now minus 1 minute) | long |
| to | End minute (default is now) | long |
| filter | An identifier for different filtering strategies. For example, filter = 001 filters out non-US clicks | long |
Table 2 Request parameters for /v1/ads/{:ad_id}/aggregated_count
Response:
| Field | Description | Type |
|---|---|---|
| ad_id | The identifier of the ad | string |
| count | The aggregated count between the start and end minutes | long |
Table 3 Response for /v1/ads/{:ad_id}/aggregated_count
API 2: Return top N most clicked ad_ids in the last M minutes
| API | Detail |
|---|---|
| GET /v1/ads/popular_ads | Return top N most clicked ads in the last M minutes |
Table 4 API for /v1/ads/popular_ads
Request parameters are:
| Field | Description | Type |
|---|---|---|
| count | Top N most clicked ads | integer |
| window | The aggregation window size (M) in minutes | integer |
| filter | An identifier for different filtering strategies | long |
Table 5 Request parameters for /v1/ads/popular_ads
Response:
| Field | Description | Type |
|---|---|---|
| ad_ids | A list of the most clicked ads | array |
Table 6 Response for /v1/ads/popular_ads
Data model
There are two types of data in the system: raw data and aggregated data.
Raw data
Below shows what the raw data looks like in log files:
[AdClickEvent] ad001, 2021-01-01 00:00:01, user 1, 207.148.22.22, USATable 7 lists what the data fields look like in a structured way. Data is scattered on different application servers.
| ad_id | click_timestamp | user | ip | country |
|---|---|---|---|---|
| ad001 | 2021-01-01 00:00:01 | user1 | 207.148.22.22 | USA |
| ad001 | 2021-01-01 00:00:02 | user1 | 207.148.22.22 | USA |
| ad002 | 2021-01-01 00:00:02 | user2 | 209.153.56.11 | USA |
Table 7 Raw data
Aggregated data
Assume that ad click events are aggregated every minute. Table 8 shows the aggregated result.
| ad_id | click_minute | count |
|---|---|---|
| ad001 | 202101010000 | 5 |
| ad001 | 202101010001 | 7 |
Table 8 Aggregated data
To support ad filtering, we add an additional field called filter_id to the table. Records with the same ad_id and click_minute are grouped by filter_id as shown in Table 9, and filters are defined in Table 10.
| ad_id | click_minute | filter_id | count |
|---|---|---|---|
| ad001 | 202101010000 | 0012 | 2 |
| ad001 | 202101010000 | 0023 | 3 |
| ad001 | 202101010001 | 0012 | 1 |
| ad001 | 202101010001 | 0023 | 6 |
Table 9 Aggregated data with filters
| filter_id | region | IP | user_id |
|---|---|---|---|
| 0012 | US | * | * |
| 0013 | * | 123.1.2.3 | * |
Table 10 Filter table
To support the query to return the top N most clicked ads in the last M minutes, the following structure is used.
| most_clicked_ads | ||
|---|---|---|
| window_size | integer | The aggregation window size (M) in minutes |
| update_time_minute | timestamp | Last updated timestamp (in 1-minute granularity) |
| most_clicked_ads | array | List of ad IDs in JSON format. |
Table 11 Support top N most clicked ads in the last M minutes
Comparison
The comparison between storing raw data and aggregated data is shown below:
| Raw data only | Aggregated data only | |
|---|---|---|
| Pros | Full data setSupport data filter and recalculation | Smaller data setFast query |
| Cons | Huge data storageSlow query | Data loss. This is derived data. For example, 10 entries might be aggregated to 1 entry |
Table 12 Raw data vs aggregated data
Should we store raw data or aggregated data? Our recommendation is to store both. Let’s take a look at why.
-
It’s a good idea to keep the raw data. If something goes wrong, we could use the raw data for debugging. If the aggregated data is corrupted due to a bad bug, we can recalculate the aggregated data from the raw data, after the bug is fixed.
-
Aggregated data should be stored as well. The data size of the raw data is huge. The large size makes querying raw data directly very inefficient. To mitigate this problem, we run read queries on aggregated data.
-
Raw data serves as backup data. We usually don’t need to query raw data unless recalculation is needed. Old raw data could be moved to cold storage to reduce costs.
-
Aggregated data serves as active data. It is tuned for query performance.
Choose the right database
When it comes to choosing the right database, we need to evaluate the following:
-
What does the data look like? Is the data relational? Is it a document or a blob?
-
Is the workflow read-heavy, write-heavy, or both?
-
Is transaction support needed?
-
Do the queries rely on many online analytical processing (OLAP) functions 3 like SUM, COUNT?
Let’s examine the raw data first. Even though we don’t need to query the raw data during normal operations, it is useful for data scientists or machine learning engineers to study user response prediction, behavioral targeting, relevance feedback, etc. 4.
As shown in the back of the envelope estimation, the average write QPS is 10,000, and the peak QPS can be 50,000, so the system is write-heavy. On the read side, raw data is used as backup and a source for recalculation, so in theory, the read volume is low.
Relational databases can do the job, but scaling the write can be challenging. NoSQL databases like Cassandra and InfluxDB are more suitable because they are optimized for write and time-range queries.
Another option is to store the data in Amazon S3 using one of the columnar data formats like ORC 5, Parquet 6, or AVRO 7. We could put a cap on the size of each file (say, 10GB) and the stream processor responsible for writing the raw data could handle the file rotation when the size cap is reached. Since this setup may be unfamiliar for many, in this design we use Cassandra as an example.
For aggregated data, it is time-series in nature and the workflow is both read and write heavy. This is because, for each ad, we need to query the database every minute to display the latest aggregation count for customers. This feature is useful for auto-refreshing the dashboard or triggering alerts in a timely manner. Since there are two million ads in total, the workflow is read-heavy. Data is aggregated and written every minute by the aggregation service, so it’s write-heavy as well. We could use the same type of database to store both raw data and aggregated data.
Now we have discussed query API design and data model, let’s put together the high-level design.
High-level design
In real-time big data 8 processing, data usually flows into and out of the processing system as unbounded data streams. The aggregation service works in the same way; the input is the raw data (unbounded data streams), and the output is the aggregated results (see Figure 2).
Asynchronous processing
The design we currently have is synchronous. This is not good because the capacity of producers and consumers is not always equal. Consider the following case; if there is a sudden increase in traffic and the number of events produced is far beyond what consumers can handle, consumers might get out-of-memory errors or experience an unexpected shutdown. If one component in the synchronous link is down, the whole system stops working.
A common solution is to adopt a message queue (Kafka) to decouple producers and consumers. This makes the whole process asynchronous and producers/consumers can be scaled independently.
Putting everything we have discussed together, we come up with the high-level design as shown in Figure 3. Log watcher, aggregation service, and database are decoupled by two message queues. The database writer polls data from the message queue, transforms the data into the database format, and writes it to the database.
What is stored in the first message queue? It contains ad click event data as shown in Table 13.
| ad_id | click_timestamp | user_id | ip | country |
|---|
Table 13 Data in the first message queue
What is stored in the second message queue? The second message queue contains two types of data:
- Ad click counts aggregated at per-minute granularity.
| ad_id | click_minute | count |
|---|
Table 14 Data in the second message queue
- Top N most clicked ads aggregated at per-minute granularity.
| update_time_minute | most_clicked_ads |
|---|
Table 15 Data in the second message queue
You might be wondering why we don’t write the aggregated results to the database directly. The short answer is that we need the second message queue like Kafka to achieve end-to-end exactly-once semantics (atomic commit).
Next, let’s dig into the details of the aggregation service.
Aggregation service
The MapReduce framework is a good option to aggregate ad click events. The directed acyclic graph (DAG) is a good model for it 9. The key to the DAG model is to break down the system into small computing units, like the Map/Aggregate/Reduce nodes, as shown in Figure 5.
Each node is responsible for one single task and it sends the processing result to its downstream nodes.
Map node
A Map node reads data from a data source, and then filters and transforms the data. For example, a Map node sends ads with ad_id % 2 = 0 to node 1, and the other ads go to node 2, as shown in Figure 6.
You might be wondering why we need the Map node. An alternative option is to set up Kafka partitions or tags and let the aggregate nodes subscribe to Kafka directly. This works, but the input data may need to be cleaned or normalized, and these operations can be done by the Map node. Another reason is that we may not have control over how data is produced and therefore events with the same ad_id might land in different Kafka partitions.
Aggregate node
An Aggregate node counts ad click events by ad_id in memory every minute. In the MapReduce paradigm, the Aggregate node is part of the Reduce. So the map-aggregate-reduce process really means map-reduce-reduce.
Reduce node
A Reduce node reduces aggregated results from all “Aggregate” nodes to the final result. For example, as shown in Figure 7, there are three aggregation nodes and each contains the top 3 most clicked ads within the node. The Reduce node reduces the total number of most clicked ads to 3.
The DAG model represents the well-known MapReduce paradigm. It is designed to take big data and use parallel distributed computing to turn big data into little- or regular-sized data.
In the DAG model, intermediate data can be stored in memory and different nodes communicate with each other through either TCP (nodes running in different processes) or shared memory (nodes running in different threads).
Main use cases
Now that we understand how MapReduce works at the high level, let’s take a look at how it can be utilized to support the main use cases:
-
Aggregate the number of clicks of ad_id in the last M mins.
-
Return top N most clicked ad_ids in the last M minutes.
-
Data filtering.
Use case 1: aggregate the number of clicks
As shown in Figure 8, input events are partitioned by ad_id (ad_id % 3) in Map nodes and are then aggregated by Aggregation nodes.
Use case 2: return top N most clicked ads
Figure 9 shows a simplified design of getting the top 3 most clicked ads, which can be extended to top N. Input events are mapped using ad_id and each Aggregate node maintains a heap data structure to get the top 3 ads within the node efficiently. In the last step, the Reduce node reduces 9 ads (top 3 from each aggregate node) to the top 3 most clicked ads every minute.
Use case 3: data filtering
To support data filtering like “show me the aggregated click count for ad001 within the USA only”, we can pre-define filtering criteria and aggregate based on them. For example, the aggregation results look like this for ad001 and ad002:
| ad_id | click_minute | country | count |
|---|---|---|---|
| ad001 | 202101010001 | USA | 100 |
| ad001 | 202101010001 | GPB | 200 |
| ad001 | 202101010001 | others | 3000 |
| ad002 | 202101010001 | USA | 10 |
| ad002 | 202101010001 | GPB | 25 |
| ad002 | 202101010001 | others | 12 |
Table 16 Aggregation results (filter by country)
This technique is called the star schema 11, which is widely used in data warehouses. The filtering fields are called dimensions. This approach has the following benefits:
-
It is simple to understand and build.
-
The current aggregation service can be reused to create more dimensions in the star schema. No additional component is needed.
-
Accessing data based on filtering criteria is fast because the result is pre-calculated.
A limitation with this approach is that it creates many more buckets and records, especially when we have a lot of filtering criteria.
Finished reading?
Mark it complete to track your progress.