Skip to content

Commit 16be614

Browse files
authored
Rjurney/motif tutorial code min (#520)
* Converted tests to pytest. Build a Python package. Update requirements.txt and split out requirements-dev.txt. Version bumps. * Restore Python .gitignore * Extra newline removed * Added VERSION file set to 0.8.5 * isort; fiex edgesDF variable name. * Back out Dockerfile changes * Back out version change in build.sbt * Backout changes to config and run-tests * Back out pytest conversion * Back out version changes to make nose tests pass * Remove changes to requirements * Put nose back in requirements.txt * Remove version bump to version.sbt * Remove packages related to testing * Remove old setup.py / setup.cfg * New pyproject.toml and poetry.lock * Short README for Python package, poetry won't allow a ../README.md path * Remove requirements files in favor of pyproject.toml * Try to poetrize CI build * pyspark min 3.4 * Local python README in pyproject.toml * Trying to remove he working folder to debug scala issue * Set Python working directory again * Accidental newline * Install Python for test... * Run tests from python/ folder * Try running tests from python/ * poetry run the unit tests * poetry run the tests * Try just using 'python' instead of a path * poetry run the last line, graphframes.main * Remove test/ folder from style paths, it doesn't exist * Remove .vscode * VERSION back to 0.8.4 * Remove tutorials reference * VERSION is a Python thing, it belongs in python/ * Include the README.md and LICENSE in the Python package * Some classifiers for pyproject.toml * Trying poetry install action instead of manual install * Removing SPARK_HOME * Returned SPARK_HOME settings * Minimized the PR to just these files * Created tutorials dependency group to minimize main bloat * Make motif.py execute in whole again * Minor isort format and cleanup of download.py * Minor isort format and cleanup of utils.py * Removed case sensitivity from the script - that was confusing people who just pasted or tried to run the code without a new SparkSession. * motif.py now matches tutorial code, runs and handles case insensitivity. * Setup a 'graphframes stackexchange' comand. * Make graphframes.tutorials.motif use a checkpoint dir unique, and from SparkSession.sparkContext. Use click.echo instead of print * Use spark.sparkContext.setCheckpointDir directly instead of instantiating a SparkContext. print-->click.echo * Using 'from __future__ import annotations' intsead of List and Tuple * Now retry three times if we can't connect for any reason in 'graphframes stackexchange' command.
1 parent fb14eff commit 16be614

8 files changed

Lines changed: 1890 additions & 4 deletions

File tree

‎python/MANIFEST.in‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,3 +7,4 @@ recursive-exclude * __pycache__
77
recursive-exclude * *.pyc
88
include README.md
99
include LICENSE
10+
include graphframes/tutorials/data/.exists

‎python/graphframes/console.py‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
import click
2+
from graphframes.tutorials import download
3+
4+
5+
@click.group()
6+
def cli():
7+
"""GraphFrames CLI: a collection of commands for graphframes."""
8+
pass
9+
10+
11+
cli.add_command(download.stackexchange)
12+
13+
14+
def main():
15+
cli()
16+
17+
18+
if __name__ == "__main__":
19+
main()
Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
1+
#!/usr/bin/env python
2+
3+
"""Download and decompress the Stack Exchange data dump from the Internet Archive."""
4+
5+
import os
6+
7+
import click
8+
import py7zr
9+
import requests # type: ignore
10+
11+
12+
@click.command()
13+
@click.argument("subdomain")
14+
@click.option(
15+
"--data-dir",
16+
default="python/graphframes/tutorials/data",
17+
help="Directory to store downloaded files",
18+
)
19+
@click.option(
20+
"--extract/--no-extract", default=True, help="Whether to extract the archive after download"
21+
)
22+
def stackexchange(subdomain: str, data_dir: str, extract: bool) -> None:
23+
"""Download Stack Exchange archive for a given SUBDOMAIN.
24+
25+
Example: python/graphframes/tutorials/download.py stats.meta
26+
27+
Note: This won't work for stackoverflow.com archives due to size.
28+
"""
29+
# Create data directory if it doesn't exist
30+
os.makedirs(data_dir, exist_ok=True)
31+
32+
# Construct archive URL and filename
33+
archive_url = f"https://archive.org/download/stackexchange/{subdomain}.stackexchange.com.7z"
34+
archive_path = os.path.join(data_dir, f"{subdomain}.stackexchange.com.7z")
35+
36+
click.echo(f"Downloading archive from {archive_url}")
37+
38+
try:
39+
# Download the file with retries
40+
max_retries = 3
41+
retry_count = 0
42+
43+
while retry_count < max_retries:
44+
try:
45+
response = requests.get(archive_url, stream=True)
46+
response.raise_for_status() # Raise exception for bad status codes
47+
break
48+
except (
49+
requests.exceptions.RequestException,
50+
requests.exceptions.ConnectionError,
51+
requests.exceptions.HTTPError,
52+
requests.exceptions.Timeout,
53+
) as e:
54+
retry_count += 1
55+
if retry_count == max_retries:
56+
click.echo(f"Failed to download after {max_retries} attempts: {e}", err=True)
57+
raise click.Abort()
58+
click.echo(f"Download attempt {retry_count} failed, retrying...")
59+
60+
total_size = int(response.headers.get("content-length", 0))
61+
62+
with click.progressbar(length=total_size, label="Downloading") as bar: # type: ignore
63+
with open(archive_path, "wb") as f:
64+
for chunk in response.iter_content(chunk_size=8192):
65+
if chunk:
66+
f.write(chunk)
67+
bar.update(len(chunk))
68+
69+
click.echo(f"Download complete: {archive_path}")
70+
71+
# Extract if requested
72+
if extract:
73+
click.echo("Extracting archive...")
74+
output_dir = f"{subdomain}.stackexchange.com"
75+
with py7zr.SevenZipFile(archive_path, mode="r") as z:
76+
z.extractall(path=os.path.join(data_dir, output_dir))
77+
click.echo(f"Extraction complete: {output_dir}")
78+
79+
except requests.exceptions.RequestException as e:
80+
click.echo(f"Error downloading archive: {e}", err=True)
81+
raise click.Abort()
82+
except py7zr.Bad7zFile as e:
83+
click.echo(f"Error extracting archive: {e}", err=True)
84+
raise click.Abort()
85+
86+
87+
if __name__ == "__main__":
88+
stackexchange()
Lines changed: 203 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,203 @@
1+
"""Demonstrate GraphFrames network motif finding capabilities. Code from the Network Motif Finding Tutorial."""
2+
3+
#
4+
# Interactive Usage: pyspark --packages graphframes:graphframes:0.8.4-spark3.5-s_2.12
5+
#
6+
# Batch Usage: spark-submit --packages graphframes:graphframes:0.8.4-spark3.5-s_2.12 python/graphframes/tutorials/motif.py
7+
#
8+
9+
import click
10+
import pyspark.sql.functions as F
11+
from pyspark.sql import DataFrame, SparkSession
12+
13+
from graphframes import GraphFrame
14+
15+
# Initialize a SparkSession
16+
spark: SparkSession = SparkSession.builder.appName("Stack Overflow Motif Analysis").getOrCreate()
17+
spark.sparkContext.setCheckpointDir("/tmp/graphframes-checkpoints/motif")
18+
19+
# Change me if you download a different stackexchange site
20+
STACKEXCHANGE_SITE = "stats.meta.stackexchange.com"
21+
BASE_PATH = f"python/graphframes/tutorials/data/{STACKEXCHANGE_SITE}"
22+
23+
24+
#
25+
# Load the nodes and edges from disk, repartition, checkpoint [plan got long for some reason] and cache.
26+
#
27+
28+
# We created these in stackexchange.py from Stack Exchange data dump XML files
29+
NODES_PATH: str = f"{BASE_PATH}/Nodes.parquet"
30+
nodes_df: DataFrame = spark.read.parquet(NODES_PATH)
31+
32+
# Repartition the nodes to give our motif searches parallelism
33+
nodes_df = nodes_df.repartition(50).checkpoint().cache()
34+
35+
# We created these in stackexchange.py from Stack Exchange data dump XML files
36+
EDGES_PATH: str = f"{BASE_PATH}/Edges.parquet"
37+
edges_df: DataFrame = spark.read.parquet(EDGES_PATH)
38+
39+
# Repartition the edges to give our motif searches parallelism
40+
edges_df = edges_df.repartition(50).checkpoint().cache()
41+
42+
# What kind of nodes we do we have to work with?
43+
node_counts = (
44+
nodes_df.select("id", F.col("Type").alias("Node Type"))
45+
.groupBy("Node Type")
46+
.count()
47+
.orderBy(F.col("count").desc())
48+
# Add a comma formatted column for display
49+
.withColumn("count", F.format_number(F.col("count"), 0))
50+
)
51+
node_counts.show()
52+
53+
# What kind of edges do we have to work with?
54+
edge_counts = (
55+
edges_df.select("src", "dst", F.col("relationship").alias("Edge Type"))
56+
.groupBy("Edge Type")
57+
.count()
58+
.orderBy(F.col("count").desc())
59+
# Add a comma formatted column for display
60+
.withColumn("count", F.format_number(F.col("count"), 0))
61+
)
62+
edge_counts.show()
63+
64+
g = GraphFrame(nodes_df, edges_df)
65+
66+
g.vertices.show(10)
67+
click.echo(f"Node columns: {g.vertices.columns}")
68+
69+
g.edges.sample(0.0001).show(10)
70+
71+
# Sanity test that all edges have valid ids
72+
edge_count = g.edges.count()
73+
valid_edge_count = (
74+
g.edges.join(g.vertices, on=g.edges.src == g.vertices.id)
75+
.select("src", "dst", "relationship")
76+
.join(g.vertices, on=g.edges.dst == g.vertices.id)
77+
.count()
78+
)
79+
80+
# Just up and die if we have edges that point to non-existent nodes
81+
assert (
82+
edge_count == valid_edge_count
83+
), f"Edge count {edge_count} != valid edge count {valid_edge_count}"
84+
click.echo(f"Edge count: {edge_count:,} == Valid edge count: {valid_edge_count:,}")
85+
86+
# G4: Continuous Triangles
87+
paths = g.find("(a)-[e1]->(b); (b)-[e2]->(c); (c)-[e3]->(a)")
88+
89+
# Show the first path
90+
paths.show(3)
91+
92+
graphlet_type_df = paths.select(
93+
F.col("a.Type").alias("A_Type"),
94+
F.col("e1.relationship").alias("(a)-[e1]->(b)"),
95+
F.col("b.Type").alias("B_Type"),
96+
F.col("e2.relationship").alias("(b)-[e2]->(c)"),
97+
F.col("c.Type").alias("C_Type"),
98+
F.col("e3.relationship").alias("(c)-[e3]->(a)"),
99+
)
100+
101+
graphlet_count_df = (
102+
graphlet_type_df.groupby(
103+
"A_Type", "(a)-[e1]->(b)", "B_Type", "(b)-[e2]->(c)", "C_Type", "(c)-[e3]->(a)"
104+
)
105+
.count()
106+
.orderBy(F.col("count").desc())
107+
# Add a comma formatted column for display
108+
.withColumn("count", F.format_number(F.col("count"), 0))
109+
)
110+
graphlet_count_df.show()
111+
112+
# G5: Divergent Triangles
113+
paths = g.find("(a)-[e1]->(b); (a)-[e2]->(c); (c)-[e3]->(b)")
114+
115+
graphlet_type_df = paths.select(
116+
F.col("a.Type").alias("A_Type"),
117+
F.col("e1.relationship").alias("(a)-[e1]->(b)"),
118+
F.col("b.Type").alias("B_Type"),
119+
F.col("e2.relationship").alias("(a)-[e2]->(c)"),
120+
F.col("c.Type").alias("C_Type"),
121+
F.col("e3.relationship").alias("(c)-[e3]->(b)"),
122+
)
123+
124+
graphlet_count_df = (
125+
graphlet_type_df.groupby(
126+
"A_Type", "(a)-[e1]->(b)", "B_Type", "(a)-[e2]->(c)", "C_Type", "(c)-[e3]->(b)"
127+
)
128+
.count()
129+
.orderBy(F.col("count").desc())
130+
# Add a comma formatted column for display
131+
.withColumn("count", F.format_number(F.col("count"), 0))
132+
)
133+
graphlet_count_df.show()
134+
135+
# G17: A directed 3-path is a surprisingly diverse graphlet
136+
paths = g.find("(a)-[e1]->(b); (b)-[e2]->(c); (d)-[e3]->(c)")
137+
138+
# Visualize the four-path by counting instances of paths by node / edge type
139+
graphlet_type_df = paths.select(
140+
F.col("a.Type").alias("A_Type"),
141+
F.col("e1.relationship").alias("(a)-[e1]->(b)"),
142+
F.col("b.Type").alias("B_Type"),
143+
F.col("e2.relationship").alias("(b)-[e2]->(c)"),
144+
F.col("c.Type").alias("C_Type"),
145+
F.col("e3.relationship").alias("(d)-[e3]->(c)"),
146+
F.col("d.Type").alias("D_Type"),
147+
)
148+
graphlet_count_df = (
149+
graphlet_type_df.groupby(
150+
"A_Type",
151+
"(a)-[e1]->(b)",
152+
"B_Type",
153+
"(b)-[e2]->(c)",
154+
"C_Type",
155+
"(d)-[e3]->(c)",
156+
"D_Type",
157+
)
158+
.count()
159+
.orderBy(F.col("count").desc())
160+
# Add a comma formatted column for display
161+
.withColumn("count", F.format_number(F.col("count"), 0))
162+
)
163+
graphlet_count_df.show()
164+
165+
graphlet_count_df.orderBy(
166+
[
167+
"A_Type",
168+
"(a)-[e1]->(b)",
169+
"B_Type",
170+
"(b)-[e2]->(c)",
171+
"C_Type",
172+
"(d)-[e3]->(c)",
173+
"D_Type",
174+
],
175+
ascending=False,
176+
).show(104)
177+
178+
# A user answers an answer that answers a question that links to an answer.
179+
linked_vote_paths = paths.filter(
180+
(F.col("a.Type") == "Vote")
181+
& (F.col("e1.relationship") == "CastFor")
182+
& (F.col("b.Type") == "Question")
183+
& (F.col("e2.relationship") == "Links")
184+
& (F.col("c.Type") == "Question")
185+
& (F.col("e3.relationship") == "CastFor")
186+
& (F.col("d.Type") == "Vote")
187+
)
188+
189+
# Sanity check the count - it should match the table above
190+
linked_vote_paths.count()
191+
192+
b_vote_counts = linked_vote_paths.select("a", "b").distinct().groupBy("b").count()
193+
c_vote_counts = linked_vote_paths.select("c", "d").distinct().groupBy("c").count()
194+
195+
linked_vote_counts = (
196+
linked_vote_paths.filter((F.col("a.VoteTypeId") == 2) & (F.col("d.VoteTypeId") == 2))
197+
.select("b", "c")
198+
.join(b_vote_counts, on="b", how="inner")
199+
.withColumnRenamed("count", "b_count")
200+
.join(c_vote_counts, on="c", how="inner")
201+
.withColumnRenamed("count", "c_count")
202+
)
203+
linked_vote_counts.stat.corr("b_count", "c_count")

0 commit comments

Comments
 (0)