Scaling Social Influence: Distributed Computation of Unified Diffusion Models

A distributed algorithm for the efficient computation of the unified model of social influence on massive datasets

2017-09-01
Alex Popa, Marc Frîncu, Charalampos Chelmis
Summary
Problem
Method
Results
Takeaways
Abstract

This paper presents a distributed, vertex-centric algorithm for the Unified Model of social influence using Apache Giraph. By mapping the complex unified diffusion formula to the Bulk Synchronous Parallel (BSP) model, the authors achieve a 3.2x performance speedup on massive datasets compared to single-node executions.

TL;DR

Researchers have developed a distributed vertex-centric algorithm to compute the Unified Model of Social Influence on massive datasets. By leveraging Apache Giraph and the Bulk Synchronous Parallel (BSP) model, the implementation achieves significant scalability, enabling the analysis of influence propagation—like rumor spreading—across networks too large for a single machine's memory.

Background: The Need for Scale

In the era of "Big Data," social networks like Facebook and Twitter generate hundreds of terabytes daily. Understanding how a "rumor" or a "viral product" spreads through these networks is critical. While previous models relied on thousands of Monte Carlo simulations (which are computationally exhausting), the Unified Model offers a direct analytical approach. However, even this model hits a wall when the graph size exceeds the RAM of a single server.

The Challenge: Memory vs. Complexity

The core problem with social graph analysis is that it is memory-intensive rather than compute-bound. Sparse graphs exhibit poor data locality, leading to frequent memory access latencies.

The Unified Model's equation for a node at time () depends on:

  • The node's previous state ().
  • Individual influence from all incoming neighbors ().
  • Collective pressure from the community ().

Storing the full history of every node's infection probability in a distributed environment would quickly crash a standard commodity cluster.

Methodology: "Thinking Like a Vertex"

The authors propose a vertex-centric approach using the Bulk Synchronous Parallel (BSP) model.

1. The Algorithm Design

Each execution is divided into "supersteps." In each step:

  1. Compute: Each vertex calculates its new infection probability using the unified formula.
  2. Communicate: Vertices send their current state to their neighbors.
  3. Synchronize: A global barrier ensures all nodes finish before the next step.

2. Memory Optimization

To prevent memory overflow, the authors optimized what is stored:

  • ICM Model: Only the last two values of are stored locally.
  • Collective History: Instead of a full list, they store the product of collective influence , significantly reducing the storage footprint.
  • Precision Scaling: Switching from double to single precision allowed processing 250,000 nodes on 4 machines where double precision would have failed.

Distributed Algorithm Pseudocode Figure 1: Conceptual mapping of node infection to the vertex-centric BSP model.

Experiments & Results

The team tested their implementation on Google Cloud Platform (GCP) using real-world datasets: the HEPT collaboration network and various sizes of the Twitter mention network.

Performance Highlights:

  • Strong Scaling: Achieved a 3.2x speedup on 8 nodes. While lower than the 5.5x seen in shared-memory (OpenMP) setups, the distributed version is capable of handling much larger total data volumes.
  • Weak Scaling: The runtime remained relatively stable as the graph size and the number of machines were doubled simultaneously, proving the algorithm's viability for "Big Data."
  • The Overhead Reality: A significant finding was that 45% of the runtime is consumed by the platform (Hadoop/Giraph) rather than the actual math.

Scaling Results Figure 2: Scaling results for the Twitter dataset demonstrating the relationship between node count and execution time.

Deep Insight: Is the Cloud Worth It?

The authors offer an honest assessment of the Cost-Performance trade-off. Because of the communication overhead of distributed systems, the speedup comes at roughly 6x the cost of a single node.

Conclusion: Distributed frameworks like Giraph are not a "silver bullet" for small-to-medium graphs where a single high-RAM server would suffice. However, for massive datasets that simply cannot fit on one machine, this distributed unified model is a necessary and effective evolution.

Future Outlook

The next frontier involves moving from "Vertex-Centric" to "Subgraph-Centric" processing. By partitioning the graph into local clusters and minimizing the messages sent across the network, researchers hope to break the current 3.2x scalability ceiling and further reduce cloud compute costs.

Find Similar Papers

Try Our Examples

  • Find recent papers that optimize Apache Giraph or Pregel-like systems specifically for memory-intensive sparse graph analysis to reduce platform overhead.
  • Which paper first proposed the "Unified Model" involving both pairwise influence and collective dynamics, and how does it mathematically differ from the Independent Cascade Model?
  • Are there any studies applying subgraph-centric frameworks like GoFFish to the Unified Model of social influence to reduce inter-node communication bottlenecks?
Contents
Scaling Social Influence: Distributed Computation of Unified Diffusion Models
1. TL;DR
2. Background: The Need for Scale
3. The Challenge: Memory vs. Complexity
4. Methodology: "Thinking Like a Vertex"
4.1. 1. The Algorithm Design
4.2. 2. Memory Optimization
5. Experiments & Results
5.1. Performance Highlights:
6. Deep Insight: Is the Cloud Worth It?
7. Future Outlook