Commit 6a0f34c
authored
feat: random walks and embeddings (#752)
* edges sampling API (scala)
* add seed to z-estimation
* wip
* WIP
* scalfix
* docstrings to RandomWalkBase and RandomWalkWithRestart
- Added Scala-style docstrings to all classes, traits, methods, and fields
- Improved documentation for random walk algorithms and configurations
* Fix RandomWalk implementation bugs and add example
- Correct element_at index from 0 to 1 for 1-based Spark SQL arrays
- Fix walk array construction by appending nextNode instead of currVisitingVertex
- Add null handling for nodes with no outgoing neighbors in restart logic
- Add comprehensive Scala docstrings to RandomWalkBase and RandomWalkWithRestart
- Create RWExample.scala demonstrating RandomWalkWithRestart on LDBC datasets
...
* Add Word2VecHashingTrick implementation for graph embeddings
This commit introduces a new Word2Vec-based embedding method using the hashing trick to handle large vocabularies efficiently in graph frames, particularly for random walk sequences. It includes configurable parameters like number of hashing functions, max features, and standard W2V settings, with comprehensive Scaladoc for public APIs.
- Added core/src/main/scala/org/graphframes/embeddings/Word2VecHashingTrick.scala: New class implementing hashing trick by applying multiple Murmur3 hash functions and modulo to map features to a fixed-size space, reducing collisions and memory usage. It trains a W2V model on expanded sequences and provides a companion model class for vector retrieval via averaging hashed embeddings. Setters include docstrings explaining trade-offs (e.g., more hashes improve quality but multiply dataset size).
- Modified core/src/main/scala/org/graphframes/examples/RWExample.scala: Updated main method to accept a single file path argument for edge loading instead of downloading LDBC datasets, simplifying usage for local files. Replaced vertex loading with direct derivation from edges for consistency and reduced I/O.
- Modified core/src/main/scala/org/graphframes/exceptions.scala: Added GraphFramesW2VException class to handle W2V-specific errors, such as unsupported input types in hashing.
* Implement reservoir sampling for neighbor selection in random walks
Replace collect_set + shuffle + slice with ReservoirSamplingAgg UDAF for
efficient sampling of up to maxNbrs neighbors per vertex. This improves
performance by avoiding full neighbor list aggregation and shuffling,
especially beneficial for high-degree vertices.
- Add ReservoirSamplingAgg trait: generic aggregator using reservoir
sampling algorithm, supporting merge operations for distributed
computation.
- Handle various vertex ID types (String, Short, Byte, Int, Long) with
appropriate encoders.
- Raise GraphFramesUnsupportedVertexTypeException for unsupported types.
- Add comprehensive test suite covering reduce, merge, and finish
operations with edge cases and fixed seeds for determinism.
Modified files:
- .gitignore: Ignore Emacs temp files for cleaner diffs.
- core/src/main/scala/org/graphframes/exceptions.scala: New exception class.
- core/src/main/scala/org/graphframes/rw/RandomWalkBase.scala: Integrate
ReservoirSamplingAgg in prepareGraph method.
New files:
- core/src/main/scala/org/apache/spark/sql/graphframes/expressions/ReservoirSamplingAgg.scala
- core/src/test/scala/org/apache/spark/sql/graphframes/expressions/ReservoirSamplingAggSuite.scala
* fix scalastyle?
* fix reservoir
* add hash2vec
delete wrong implementation of w2v + hashing
* fixes
* docstrings + scalfix
* remove sampling as not needed
* fixes in build and code
* Big update
- replace Reservoir sampling by KMinSampling
- add L2norm to Hash2vec
- add an optional convolution step to RW embeddings
- small updates and performance fixes
* Tests and updates
* workaround scala 2.13 deprecation of Searching.search
* Fixes
Fix the problem `sun.security.action` access
* Fix access
* fallback to java Serialization
* Fix some bugs
Tested on an end2end case
* Sampling Convolution tests and docstrings
* docstrings for RW Emebddings
* Python API
* hash2vec tests
* Explicit types
* hash2vec and random walks with restart tests
* performance
* ignore unused nowarn
spark 4 vs spark 3
* performance + cached walks support + continous mode
* fix rw and update the branch
* fix
* protobuf & connect
* public API for embeddings and small refactoring of methods
* initial Py API for embeddings
* fix
* tests + docs
- python tests
- docs (I use AI to generate but checked by myself and fixed)
- fix in python APIs
- small changes
* fix connect tests
* decrease GC pressure
* further optimizations
- reduce GC pressure from UTF8String
- avoid hash recomputations
* refactor: optimize Hash2Vec string hashing and partitioning logic
* refactor: replace generic hash function with type-specific implementations and add hash caching
* refactor: simplify hash function logic and improve performance in Hash2Vec
* refactor: inline generic processPartitionGeneric into specialized String and Long methods
* test: add tests for PagedMatrixDouble helper covering page extension, add, and getVector
* refactor: replace unsafe hash functions with MurmurHash3 and optimize memory usage with paged matrix
* refactor: update processStringPartition to use PagedMatrixDouble for memory efficiency
* docs: add internal documentation for PagedMatrixDouble explaining memory layout and GC benefits
* chore: remove commented helper section from Hash2Vec
* refactor: replace case class with class for PagedMatrixDouble and remove Scala version-specific compat files
* fix: correct LongMap type parameter in Hash2Vec vocabIndex initialization
* refactor: reduce PAGE_BITS from 16 to 12 and update related constants and tests
* feat: add max vectors per partition limit and batched processing for memory control
* refactor: process long partitions in batches respecting max vectors limit
* test: add Hash2Vec tests for co-occurrence patterns and cosine similarity validation
* chore: clean up code formatting and improve test readability in Hash2Vec
* fix: correct typo in error message from 'gor' to 'got' in Hash2Vec exception
* docs: add docstrings for Hash2Vec setters setDoNormalization and setMaxVectorsPerPartition
* fix: correct typo in NOTICE and KMinSampling, update Hash2Vec defaults, add null check and memory management in RandomWalkEmbeddings, simplify RandomWalkBase seed setting, and fix randomness in RandomWalkWithRestart
* fix: skip seeds for previous batches to maintain consistency when starting from non-first batch
* fix: add overwrite mode when writing batch results to allow re-running with same walkID
* feat: add cleanUp method to remove temporary files for a walk ID using Hadoop FS
* docs: improve documentation for cleanUp method in RandomWalkBase trait
* test: add cleanUp call in RandomWalkWithRestart test
* test: move walks execution inside try block in RandomWalkWithRestartSuite
* refactor: set walkID default to UUID and remove redundant runID variable
* docs: update comment to clarify walkID retrieval method behavior
* refactor: rename walkID to runID for clarity in random walk operations
* refactor: remove runId parameter from cleanUp method in RandomWalkBase
* refactor: move cleanUp method to companion object with parameters and update instance method
* refactor: improve documentation formatting and remove redundant log setup in cleanUp
* test: verify temporary files are deleted after RandomWalkWithRestart test
* refactor: use numBatches variable and fix run path in RandomWalkWithRestartSuite test
* test: add test for RandomWalkWithRestart resuming from middle iteration
* style: Format RandomWalkWithRestartSuite with consistent spacing and indentation
* refactor: use getSeq instead of getAs[Seq[String]] in RandomWalkWithRestartSuite
* fix: correct typo in error message and parameter name for gaussian sigma
* feat: add clean-up option for temporary random walk files in embeddings pipeline
* docs: improve formatting and consistency in graph-ml documentation tables and sections
* feat: add clean_up_after_run field to RandomWalkEmbeddings proto definition
* chore: add clean_up_after_run parameter to _RandomWalksEmbeddingsParameters1 parent 12ad0c8 commit 6a0f34c
32 files changed
Lines changed: 3792 additions & 163 deletions
File tree
- connect/src/main
- protobuf
- scala/org/apache/spark/sql/graphframes
- core/src
- main/scala/org
- apache/spark/sql/graphframes/expressions
- graphframes
- convolutions
- embeddings
- examples
- rw
- test/scala/org
- apache/spark/sql/graphframes/expressions
- graphframes
- convolutions
- embeddings
- rw
- docs/src/04-user-guide
- project
- python
- graphframes
- classic
- connect
- proto
- internal
- tests
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
79 | 79 | | |
80 | 80 | | |
81 | 81 | | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
21 | 21 | | |
22 | 22 | | |
23 | 23 | | |
24 | | - | |
| 24 | + | |
25 | 25 | | |
26 | 26 | | |
27 | 27 | | |
28 | 28 | | |
29 | 29 | | |
30 | 30 | | |
31 | | - | |
| 31 | + | |
32 | 32 | | |
33 | 33 | | |
34 | 34 | | |
35 | | - | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
8 | 8 | | |
9 | 9 | | |
10 | 10 | | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
96 | 96 | | |
97 | 97 | | |
98 | 98 | | |
99 | | - | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
100 | 102 | | |
101 | 103 | | |
102 | | - | |
| 104 | + | |
103 | 105 | | |
104 | 106 | | |
105 | 107 | | |
| |||
111 | 113 | | |
112 | 114 | | |
113 | 115 | | |
114 | | - | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
115 | 120 | | |
116 | 121 | | |
117 | 122 | | |
| |||
122 | 127 | | |
123 | 128 | | |
124 | 129 | | |
125 | | - | |
126 | | - | |
| 130 | + | |
127 | 131 | | |
128 | 132 | | |
129 | 133 | | |
| |||
136 | 140 | | |
137 | 141 | | |
138 | 142 | | |
139 | | - | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
140 | 146 | | |
141 | 147 | | |
142 | 148 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
35 | 35 | | |
36 | 36 | | |
37 | 37 | | |
38 | | - | |
39 | 38 | | |
| 39 | + | |
| 40 | + | |
40 | 41 | | |
41 | 42 | | |
42 | 43 | | |
| |||
208 | 209 | | |
209 | 210 | | |
210 | 211 | | |
| 212 | + | |
| 213 | + | |
| 214 | + | |
| 215 | + | |
| 216 | + | |
| 217 | + | |
| 218 | + | |
| 219 | + | |
| 220 | + | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
| 235 | + | |
| 236 | + | |
| 237 | + | |
| 238 | + | |
| 239 | + | |
| 240 | + | |
| 241 | + | |
| 242 | + | |
| 243 | + | |
| 244 | + | |
| 245 | + | |
| 246 | + | |
Lines changed: 39 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
12 | 12 | | |
13 | 13 | | |
14 | 14 | | |
| 15 | + | |
15 | 16 | | |
16 | 17 | | |
17 | 18 | | |
| |||
454 | 455 | | |
455 | 456 | | |
456 | 457 | | |
| 458 | + | |
| 459 | + | |
| 460 | + | |
| 461 | + | |
| 462 | + | |
| 463 | + | |
| 464 | + | |
| 465 | + | |
| 466 | + | |
| 467 | + | |
| 468 | + | |
| 469 | + | |
| 470 | + | |
| 471 | + | |
| 472 | + | |
| 473 | + | |
| 474 | + | |
| 475 | + | |
| 476 | + | |
| 477 | + | |
| 478 | + | |
| 479 | + | |
| 480 | + | |
| 481 | + | |
| 482 | + | |
| 483 | + | |
| 484 | + | |
| 485 | + | |
| 486 | + | |
| 487 | + | |
| 488 | + | |
| 489 | + | |
| 490 | + | |
| 491 | + | |
| 492 | + | |
| 493 | + | |
| 494 | + | |
| 495 | + | |
457 | 496 | | |
458 | 497 | | |
459 | 498 | | |
| |||
Lines changed: 165 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
| 149 | + | |
| 150 | + | |
| 151 | + | |
| 152 | + | |
| 153 | + | |
| 154 | + | |
| 155 | + | |
| 156 | + | |
| 157 | + | |
| 158 | + | |
| 159 | + | |
| 160 | + | |
| 161 | + | |
| 162 | + | |
| 163 | + | |
| 164 | + | |
| 165 | + | |
0 commit comments