Scaling Social Discovery: Two-Tier Parallel Crawling for Online Social Networks
Parallel crawling for online social networks
This paper presents a robust parallel crawling framework specifically designed for Online Social Networks (OSNs), employing a dynamic assignment architecture via a centralized database-backed queue. The authors implemented this system to crawl eBay, successfully identifying over 11 million users and fully indexing 66,130 profiles with zero redundancy.
TL;DR
Researchers at Carnegie Mellon University developed a high-efficiency parallel crawler tailored for the structured world of Online Social Networks (OSNs). By utilizing a centralized MySQL queue to manage "agent" nodes, the system achieved zero redundancy and high fault tolerance, successfully mapping a network of over 11 million users on eBay.
Background & Positioning
In the mid-2000s, as networks like eBay, LinkedIn, and Friendster exploded in size, the academic community faced a challenge: how to gather social graph data efficiently. While general search engines had proprietary crawlers, scientific documentation on how to build these systems was scarce. This paper fills that gap by implementing a Dynamic Assignment Architecture, a model previously proposed in theory but rarely documented in practical social network contexts.
Problem & Motivation
Standard web crawling approaches often fall into two traps:
- Redundancy: In an "Independent Architecture," crawlers don't talk to each other, leading to the same expensive pages being downloaded multiple times.
- Inflexibility: "Static Assignment" (partitioning the web by URL hash) makes it difficult to add new machines mid-crawl or prioritize specific high-value users.
Social networks are unique because the data is highly structured (user IDs, feedback loops, friend lists). The authors realized that by treating the crawl as a graph traversal (BFS) rather than a random walk, they could use a central coordinator to solve both redundancy and flexibility issues.
Methodology: The Two-Tier Architecture
The core innovation lies in the Centralized Queue via MySQL. Instead of a simple in-memory list, the master uses a database to track every user ID found.
The Two Tiers of Parallelism:
- Level 1 (Agent-Level): Multiple physical machines (Agents) request work from the Master. If one machine crashes, the others continue, and the "unprocessed" task simply remains in the database.
- Level 2 (Thread-Level): Each Agent runs multiple concurrent threads, maximizing the utilization of the local network interface.

The logic follows a Breadth-First Search (BFS):
- Start with a "seed" set of users.
- Pop user IDs from the central queue.
- Crawl the user’s "Feedback" page to find links to other users.
- Insert new IDs into the queue if they haven't been seen before.
Experiments & Real-World Results
The authors tested their system on eBay, which they argue is one of the world's largest social networks due to the transaction links between buyers and sellers.

- Duration: Oct 10 to Nov 2 (approx. 23 days).
- Scale: 11,716,588 users discovered; 66,130 users fully parsed.
- Efficiency: The centralized database proved remarkably scalable. The bottleneck was not the database lock contention but the raw Internet bandwidth of the agents.
Key Insight: Manual Override
Because the system uses a centralized database, the researchers could manually "jump the queue." By inserting a specific user ID at the start of the Queue table, they could force the next available crawler thread to prioritize that user—a feature essential for monitoring real-time fraudulent behavior.
Critical Analysis & Conclusion
Takeaways
- Structured Data is an Advantage: Using unique identifiers (UIDs) instead of raw URLs allows for more precise tracking and deduplication.
- Simplicity Scales: A standard RDBMS (MySQL) is sufficient to coordinate millions of tasks, acting as a "Source of Truth" that simplifies the distributed systems engineering significantly.
Limitations & Future Work
The 2007-era paper focused on a "polite" but aggressive BFS. Modern social networks now use heavy rate-limiting and anti-bot measures that would require more sophisticated proxy rotation and "human-like" behavior simulation not covered in this framework. Furthermore, as the discovered nodes reached 11 million but only 66,000 were fully crawled, there is a clear trade-off between "discovery depth" and "data breadth" that warrants more algorithmic optimization.
This work remains a foundational example of how to bridge software engineering with graph theory to map the digital social connections that define our modern internet.
