NOTE

Handling Massive Data

A historical note on divide-and-conquer, data structures, and common massive-data interview problems.

System DesignCreated Updated 6 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is Massive Data?

The data volume is so large that it cannot be loaded into memory all at once for processing.

2. How to Process Massive Data

Divide and conquer + process with an appropriate data structure [key point] + merge the results. Since the data cannot be loaded all at once, split it and load it multiple times, use a suitable data structure for processing, and finally merge the results.

2.1. How to Split

Read a large file line by line and split it into small files using hash + modulo.

2.2. How to Choose a Data Structure

Think about how you would solve the problem if the data volume were small.

  1. Bloom Filter
    Used to quickly determine whether a value is in a set, with a certain false-positive rate.
    Suitable for deduplication and set intersection.
    BloomFilter.md

  2. Hash
    Maps a value of arbitrary length to a fixed-length integer value.
    Suitable for partitioning data, building a HashMap to count occurrences, using a HashSet to detect duplicates, etc.

  3. Bit Map
    Use one bit to represent one number: the first bit represents 1, the second bit represents 2, and so on, greatly saving space.
    Suitable for sorting, deduplication, etc.

  4. Heap
    Represent a binary tree with an array.
    Heap

    Suitable for extracting Top K numbers.

  5. Trie
    Use a tree to store common prefixes.

    Suitable for word-frequency statistics.

  6. External Sorting
    Split a large file into small files that fit in memory, sort the small files separately, and finally merge them.
    Suitable for sorting large files.

2.3. How to Merge the Results

Merge algorithm.

3. Examples

3.1. Find the Most Frequently Accessed IP [Hash-Split Small Files + HashMap]

Given massive log data, extract the IP that accessed Baidu most frequently on a certain day.

  • Analysis
    • First consider the small-data scenario. For Top 1, use HashMap<IP, count> to count occurrences, then sort and find the largest count. Top 1 can use selection sort.
    • The problem is that IPv4 has 2^32 = 4G possible IPs, so creating a 4G-sized HashMap is unrealistic and the data needs to be split.
  • Solution
    1. Read the file line by line and distribute records into 1024 small files using Hash(IP) % 1024. The HashMap for each file therefore occupies at most 4G / 1024 = 4M, and the same IP must go into the same file. In the extreme case where all IPs are identical, there will be only one resulting file, which is still fine because the input is read line by line.
    2. Read each small file line by line, count IPs with a HashMap, and save the (IP, count) with the largest count into a new HashMap.
    3. Sort this HashMap by count and find the IP with the largest count.

A search engine records every query string used by users in log files. Each query is 1-255 bytes long. Suppose there are ten million records. The queries have a high duplication rate: there are no more than three million distinct queries. The more often a query is repeated, the more popular it is. Find the 10 most popular queries using no more than 1 GB of memory.

  • Analysis
    • First consider the small-data scenario. For Top K, use HashMap<query, count> to count occurrences and then select the K largest. Top K can be handled with a heap.
    • There are no more than three million distinct queries. A HashMap would consume approximately 3,000,000 * 255 / 1024 / 1024 = 729.56M, which meets the memory limit.
  • Solution
    1. Read the file line by line and count keywords with a HashMap.
    2. Use a min-heap on the HashMap to obtain Top K.

3.3. Find the 100 Most Frequent Words [HashMap + Heap]

There is a 1 GB file, each line contains one word, and each word is no more than 16 bytes. The memory limit is 1 MB. Return the 100 most frequent words.

  • Analysis
    • First consider the small-data scenario. For Top K, use HashMap<query, count> and then obtain the largest K counts with a heap.
    • The file is 1 GB while the memory limit is 1 MB, so it first needs to be split into small files. Hash partitioning can be used.
  • Solution
    1. Read the file line by line and distribute each word into 5000 small files according to Hash(word) % 5000.
    2. Read each small file from step 1, count occurrences with a HashMap, and store the results in another HashMap.
    3. Use a min-heap of length 100 on the HashMap from step 2 to obtain Top K.
  • Similar problem

    There are 10 files, each 1 GB. Every line in every file contains a user query, and queries may repeat across files. Sort the queries by frequency.

3.4. Find Common URLs [Bloom Filter / Hash-Split Small Files + HashSet]

Given files a and b, each containing 5 billion URLs, with each URL occupying 64 bytes and a 4 GB memory limit, find the URLs common to both files.

  • Analysis
    • First consider the small-data scenario. Put URLs from file a into a HashSet, then traverse URLs from file b and check whether each URL is in the HashSet.
    • The memory limit is 4 GB, while 5 billion * 64B / 1024 / 1024 clearly exceeds the limit, so the data needs to be split.
    • If errors are acceptable, use a Bloom Filter. Five billion bits require about 5,000,000,000 / 1024 / 1024 / 1024 / 8 = 0.582GB.
  • Solution
    • Option 1: Bloom Filter
      1. Read file a and map each URL into the Bloom Filter.
      2. Read file b and test whether each URL is in the Bloom Filter. If it is, treat it as a common URL and save it.
    • Option 2: Divide and conquer
      1. Read files a and b and split them with Hash(URL) % 1000 into a001...a999 and b001...b999.
      2. Read a001 and b001, use a HashSet to test for common values, and save matches.
      3. Continue as in step 2 for the remaining files.

3.5. Find Non-Duplicate Integers [BitMap / Hash-Split Small Files + HashSet]

Find the integers that occur only once among 250 million integers. Memory is insufficient to hold all 250 million integers.

  • Analysis
    • First consider the small-data scenario. Use HashMap<number, count> and find entries whose count is 1.
    • Since memory cannot hold all 250 million integers, the data needs to be split.
    • Alternatively, use a BitMap to represent the 250 million integers, using 2 bits per number rather than 1 because duplicates need to be distinguished. This requires 5,000,000,000 * 2 / 8 / 1024 / 1024 = 1200M of memory.
  • Solution
    • Option 1
      • Read integers one by one and hash-partition them into small files.
      • Read each small file, use a HashMap to count occurrences, and store numbers with count 1 in a separate file.
      • The final file contains all numbers that occur exactly once.
    • Option 2
      • Construct a 1200 MB bitmap.
      • Read the numbers. Set the state to 01 on the first occurrence, 10 on the second or later occurrence, and 00 if absent.
      • Read out the numbers whose state is 01.

3.6. Determine Whether a Number Is in a Set [BitMap/BloomFilter]

Given 4 billion distinct, unsorted unsigned int values and then another number, how can you quickly determine whether that number is among the 4 billion values?

  • Analysis
    • First consider the small-data scenario. Store all integers in a HashSet and check whether the number is in the HashSet.
    • Memory cannot hold all 4 billion integers, so splitting is needed, but after splitting it no longer satisfies the fast-lookup requirement.
    • A BitMap can represent the 4 billion numbers and requires 4,000,000,000 / 8 / 1024 / 1024 = 476.84M.
  • Solution
    1. Use about 500 MB to represent these numbers.
    2. Read the numbers and store them in the bitmap.
    3. Read the target number and check whether the corresponding bitmap bit is 1.

3.7. Given X Numbers, Extract the Largest Y

  • If X is small, load all values into memory, sort them, and take the first Y. For example, quicksort is O(X log X).
  • If X is large and cannot be loaded into memory all at once, there is no need to load all values if only the largest Y are needed. Use a min-heap for Top K: if the heap has fewer than 100 elements, insert directly; otherwise, if a new number is greater than the heap top, remove the heap top, adjust the heap, insert the new element, and adjust again. Complexity is O(X log Y).

4. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub