Commit 8471b1a
feat(pregel): automatically skip second join when dst columns not needed (#795)
* feat(pregel): automatically skip second join when dst columns not needed
Implements automatic optimization for Pregel triplet generation that
skips the second join (adding destination vertex state) when no message
expressions reference dst.* columns.
The optimization works by:
1. Analyzing all message expressions before the iteration loop
2. Extracting column prefixes (src, dst, edge) from the expression AST
3. Skipping the dst vertex join if no dst.* columns are referenced
AND skipMessagesFromNonActiveVertices is disabled
This provides significant performance improvement for algorithms like
PageRank, directed LabelPropagation, and DetectingCycles that only
need source vertex or edge columns in their message expressions.
Closes #790
* fix: only analyze message expressions for dst detection, provide dst.id when skipping join
The previous implementation incorrectly checked both the target ID expression
and message expression for dst.* references. Since sendMsgToDst uses
Pregel.dst(ID) as the target, it would always detect 'dst' as referenced
even when the message itself only used src columns.
This fix:
1. Only analyzes the message expressions (not target ID) for dst.* references
2. When skipping the join, creates a minimal dst struct with just the id
from edge_dst so that sendMsgToDst can still route messages correctly
Added test: 'sendMsgToDst with only src columns in message' to verify
the optimization works correctly when dst.id is implicitly used for routing.
* style: run scalafmt and update SparkShims docs to be implementation-agnostic
* feat: skip dst join when only dst.id is referenced
- Add extractColumnReferences to SparkShims returning Map[String, Set[String]]
to track which specific fields are accessed under each prefix
- Handle resolved expressions (AttributeReference, GetStructField) in addition
to unresolved ones for more robust column detection
- Update Pregel optimization to skip dst join when only dst.id is referenced
since dst.id is available from the edge's dst column
- Change optimization log message from logInfo to logDebug
- Add test for dst.id-only reference case
* refactor: address reviewer feedback on dst join optimization
- Remove unused extractColumnPrefixes method from SparkShims (both Spark 3/4)
- Refactor Pregel.scala to parse expressions once instead of twice
- Add documentation for deeply nested struct access fallback behavior
* test: add comprehensive tests for extractColumnReferences and dst join optimization
- Add SparkShimsSuite with 22 unit tests for column reference extraction
- Add 4 integration tests to PregelSuite for complex dst usage patterns
- Fix UTF8String handling in UnresolvedExtractValue pattern matching
* refactor: remove Pregel references from SparkShims comments
Address peer feedback to keep SparkShims implementation-agnostic by removing
specific algorithm references from comments while maintaining functional clarity.
🤖 Generated with [Claude Code](https://claude.ai/code)
Co-Authored-By: Claude <noreply@anthropic.com>
* perf: optimize caching and partitioning when skipping dst join
Implement peer feedback suggestions:
1. Move cache/checkpoint logic before expensive operations - persist srcWithEdges
when skipping dst join to avoid recomputation
2. Change partitioning to src-only when dst join is skipped since dst
partitioning is unnecessary
3. Move dst state detection earlier to enable these optimizations
These changes provide additional performance improvements for algorithms like
PageRank that only need source vertex data.
🤖 Generated with [Claude Code](https://claude.ai/code)
Co-Authored-By: Claude <noreply@anthropic.com>
* fix: correct repartition syntax for conditional partitioning
Fix Scala syntax error in repartition call that was causing CI build failures.
Use proper sequence expansion syntax for multiple column repartitioning.
🤖 Generated with [Claude Code](https://claude.ai/code)
Co-Authored-By: Claude <noreply@anthropic.com>
* fix: correct repartition syntax for conditional partitioning
Fix Scala syntax error in repartition call that was causing CI build failures.
Use proper sequence expansion syntax for multiple column repartitioning.
🤖 Generated with [Claude Code](https://claude.ai/code)
Co-Authored-By: Claude <noreply@anthropic.com>
* style: apply scalafmt formatting to Pregel.scala
Fix formatting issues that were causing CI scalafmt checks to fail.
🤖 Generated with [Claude Code](https://claude.ai/code)
Co-Authored-By: Claude <noreply@anthropic.com>
* address Sem's PR comments
* push triplet filtering when no dst join
---------
Co-authored-by: James <james@goivio.com>
Co-authored-by: Claude <noreply@anthropic.com>1 parent f67ba04 commit 8471b1a
5 files changed
Lines changed: 690 additions & 7 deletions
File tree
- core/src
- main
- scala-spark-3/org/apache/spark/sql/graphframes
- scala-spark-4/org/apache/spark/sql/graphframes
- scala/org/graphframes/lib
- test/scala/org/graphframes
- lib
Lines changed: 73 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
22 | 22 | | |
23 | 23 | | |
24 | 24 | | |
| 25 | + | |
| 26 | + | |
25 | 27 | | |
| 28 | + | |
| 29 | + | |
26 | 30 | | |
27 | 31 | | |
28 | 32 | | |
| 33 | + | |
29 | 34 | | |
30 | 35 | | |
31 | 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 | + | |
32 | 105 | | |
33 | 106 | | |
34 | 107 | | |
| |||
Lines changed: 74 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
21 | 21 | | |
22 | 22 | | |
23 | 23 | | |
| 24 | + | |
| 25 | + | |
24 | 26 | | |
| 27 | + | |
| 28 | + | |
25 | 29 | | |
26 | 30 | | |
27 | 31 | | |
28 | 32 | | |
29 | 33 | | |
30 | 34 | | |
31 | 35 | | |
| 36 | + | |
| 37 | + | |
32 | 38 | | |
33 | 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 | + | |
34 | 108 | | |
35 | 109 | | |
36 | 110 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
| 27 | + | |
27 | 28 | | |
28 | 29 | | |
29 | 30 | | |
| |||
397 | 398 | | |
398 | 399 | | |
399 | 400 | | |
| 401 | + | |
| 402 | + | |
| 403 | + | |
| 404 | + | |
| 405 | + | |
| 406 | + | |
| 407 | + | |
| 408 | + | |
| 409 | + | |
| 410 | + | |
| 411 | + | |
| 412 | + | |
| 413 | + | |
| 414 | + | |
| 415 | + | |
| 416 | + | |
| 417 | + | |
| 418 | + | |
| 419 | + | |
| 420 | + | |
| 421 | + | |
400 | 422 | | |
401 | 423 | | |
402 | | - | |
| 424 | + | |
403 | 425 | | |
404 | 426 | | |
405 | 427 | | |
| |||
431 | 453 | | |
432 | 454 | | |
433 | 455 | | |
434 | | - | |
| 456 | + | |
| 457 | + | |
| 458 | + | |
| 459 | + | |
| 460 | + | |
| 461 | + | |
| 462 | + | |
| 463 | + | |
| 464 | + | |
435 | 465 | | |
436 | 466 | | |
437 | | - | |
438 | | - | |
439 | | - | |
440 | | - | |
441 | 467 | | |
442 | | - | |
| 468 | + | |
| 469 | + | |
| 470 | + | |
| 471 | + | |
| 472 | + | |
| 473 | + | |
| 474 | + | |
| 475 | + | |
| 476 | + | |
| 477 | + | |
| 478 | + | |
| 479 | + | |
| 480 | + | |
| 481 | + | |
| 482 | + | |
| 483 | + | |
| 484 | + | |
443 | 485 | | |
444 | 486 | | |
445 | 487 | | |
| |||
0 commit comments