Towards online graph processing with spark streaming
Bibliographic record
Abstract
Graph processing is one of the most important topics in big data processing. The graph architecture is suitable for distributed processing as the processing works in an iterative manner allowing parallelism. Also, the structure has proved to be suitable in representing social networks, web page indexes, and many other problems. However, graph processing introduce many problems as well. Partitioning the graph to distribute the data on multiple machines and minimizing data movement is a serious challenge. Also many of the graph algorithms have high complexity. GraphX is one of the frameworks that introduce an abstraction on top of Spark, an iterative data processing engine. However, GraphX and other novel graph abstractions still do not support processing data streams with online graphs. In this work we try to use IndexedRDD, a library to enable fine grained updates as a key-value store on top of Spark to represent a graph structure and test if it can be used as an efficient online graph storage for spark streaming. We did experiments to compare our data streaming implementation using IndexedRDD with the obvious elementary solution of using RDD transformations to join the old RDD with the new one to make a new composite RDD on each micro-batch. We also want to compare the above two with a distributed in-memory key-value store (such as Redis). The results show big advantage of using Redis over RDD transformations and IndexedRDD. However, it has some limitations such as lacking the support for property graphs. IndexedRDD, on the other hand, has shown good performance for insertions and a shortcoming in its need to rebuild the index after each data update, which add extra time on each lookup that cannot be tolerated when lookup speed is essential.
Fetched live from OpenAlex and de-inverted. Abstracts are not stored in this database: the inverted indexes are 8.6 GB of the frame’s 9.3 GB of text, and the host has 13 GB free.
How this classification was reachedexpand
Full frame machine prediction
Teacher imitationNot calibrated prevalence, not ground truth. Human validation pending. The Gemma side is a direct model label for every work in the frame, read from the title-only record. The Codex side is a classifier learned from the 10,348 direct Codex labels and calibrated to design-weighted sample rates; fields without enough sample support carry no Codex call. Candidate is the union of the two sides; consensus is their intersection. These outputs are machine_predicted_unvalidated and are not human labels.
Distilled classifier scores by category (both heads)
| Category | Codex | Gemma |
|---|---|---|
| Metaresearch | 0.002 | 0.009 |
| Meta-epidemiology (narrow) | 0.001 | 0.001 |
| Meta-epidemiology (broad) | 0.001 | 0.002 |
| Bibliometrics | 0.002 | 0.002 |
| Science and technology studies | 0.001 | 0.001 |
| Scholarly communication | 0.003 | 0.006 |
| Open science | 0.004 | 0.004 |
| Research integrity | 0.001 | 0.003 |
| Insufficient payload (model declined to judge) | 0.005 | 0.004 |
Machine scores (provisional)
The two teacher heads of the student model, read on this work. A score orders the frame for review; it never asserts a category, and the validation status ships verbatim with every row.
Baseline scores from an immature model (maturity gate not passed, 7 training rounds). Scores rank; they never assert a category.
score_only:v0-immature-baseline · verbatim from the scoring run: score_only means the number may rank works, and no category label ships from itClassification
machine, unvalidatedMachine predicted; a candidate call from one source (direct Gemma or distilled Codex), not a consensus.
How this classification was reached, model by model and score by score, is at the end of the page under "How this classification was reached".