NOTE
Handling Massive Data
A historical note on divide-and-conquer, data structures, and common massive-data interview problems.
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.
-
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 -
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. -
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. -
Heap
Represent a binary tree with an array.
Heap

Suitable for extracting Top K numbers. -
Trie
Use a tree to store common prefixes.

Suitable for word-frequency statistics. -
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 = 4Gpossible IPs, so creating a 4G-sized HashMap is unrealistic and the data needs to be split.
- First consider the small-data scenario. For Top 1, use
- Solution
- 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 most4G / 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. - 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. - Sort this HashMap by count and find the IP with the largest count.
- Read the file line by line and distribute records into 1024 small files using
3.2. Find the 10 Most Popular Search Keywords [HashMap + Heap]
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.
- First consider the small-data scenario. For Top K, use
- Solution
- Read the file line by line and count keywords with a HashMap.
- 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.
- First consider the small-data scenario. For Top K, use
- Solution
- Read the file line by line and distribute each word into 5000 small files according to
Hash(word) % 5000. - Read each small file from step 1, count occurrences with a HashMap, and store the results in another HashMap.
- Use a min-heap of length 100 on the HashMap from step 2 to obtain Top K.
- Read the file line by line and distribute each word into 5000 small files according to
- 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 / 1024clearly 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
- Read file a and map each URL into the Bloom Filter.
- 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
- Read files a and b and split them with
Hash(URL) % 1000intoa001...a999andb001...b999. - Read
a001andb001, use a HashSet to test for common values, and save matches. - Continue as in step 2 for the remaining files.
- Read files a and b and split them with
- Option 1: Bloom Filter
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 = 1200Mof memory.
- First consider the small-data scenario. Use
- 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
01on the first occurrence,10on the second or later occurrence, and00if absent. - Read out the numbers whose state is
01.
- Option 1
3.6. Determine Whether a Number Is in a Set [BitMap/BloomFilter]
Given 4 billion distinct, unsorted
unsigned intvalues 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
- Use about 500 MB to represent these numbers.
- Read the numbers and store them in the bitmap.
- 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
- Ten Massive-Data Interview Problems and Ten Methods
- Massive-Data Processing for Interviews
- How to Quickly Solve 99% of Massive-Data Interview Problems
- Find the Largest 10,000 Numbers Among One Billion Numbers
- Splitting and Merging Large Folders - Java
- Find the Person Who Accessed a Site Most Frequently
- External Sorting: Sort Two Billion Integers with 2 GB of Memory
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub