Skip to content
Merged
Changes from 1 commit
Commits
Show all changes
64 commits
Select commit Hold shift + click to select a range
f4e9cdb
Converted tests to pytest. Build a Python package. Update requirement…
rjurney Feb 16, 2025
c256244
Restore Python .gitignore
rjurney Feb 16, 2025
6c3df0b
Extra newline removed
rjurney Feb 16, 2025
b2838d2
Merge branch 'master' of github.com:graphframes/graphframes into rjur…
rjurney Feb 16, 2025
caf5091
Added VERSION file set to 0.8.5
rjurney Feb 16, 2025
7cfa2d1
isort; fiex edgesDF variable name.
rjurney Feb 16, 2025
2ca9a15
Merge branch 'master' of github.com:graphframes/graphframes into rjur…
rjurney Feb 16, 2025
a8bf0be
Back out Dockerfile changes
rjurney Feb 16, 2025
54a942d
Back out version change in build.sbt
rjurney Feb 16, 2025
8b0e346
Backout changes to config and run-tests
rjurney Feb 16, 2025
46c2b93
Back out pytest conversion
rjurney Feb 16, 2025
18b5da0
Back out version changes to make nose tests pass
rjurney Feb 16, 2025
8eca097
Remove changes to requirements
rjurney Feb 16, 2025
277c06f
Put nose back in requirements.txt
rjurney Feb 16, 2025
b55ee48
Remove version bump to version.sbt
rjurney Feb 16, 2025
f8a8fd9
Remove packages related to testing
rjurney Feb 16, 2025
bc2cb36
Remove old setup.py / setup.cfg
rjurney Feb 16, 2025
728be33
New pyproject.toml and poetry.lock
rjurney Feb 16, 2025
3cea1a8
Short README for Python package, poetry won't allow a ../README.md path
rjurney Feb 16, 2025
87cc975
Remove requirements files in favor of pyproject.toml
rjurney Feb 16, 2025
6f84a5a
Try to poetrize CI build
rjurney Feb 16, 2025
9a8eef0
pyspark min 3.4
rjurney Feb 16, 2025
75ecd99
Local python README in pyproject.toml
rjurney Feb 16, 2025
80231d0
Trying to remove he working folder to debug scala issue
rjurney Feb 16, 2025
2a9170b
Set Python working directory again
rjurney Feb 16, 2025
3de2263
Accidental newline
rjurney Feb 16, 2025
4662717
Install Python for test...
rjurney Feb 17, 2025
1b7b9f8
Run tests from python/ folder
rjurney Feb 17, 2025
58da493
Try running tests from python/
rjurney Feb 17, 2025
9f4aa24
poetry run the unit tests
rjurney Feb 17, 2025
11b2782
poetry run the tests
rjurney Feb 17, 2025
9772344
Try just using 'python' instead of a path
rjurney Feb 17, 2025
d55dbfe
poetry run the last line, graphframes.main
rjurney Feb 17, 2025
2fc4d08
Remove test/ folder from style paths, it doesn't exist
rjurney Feb 17, 2025
8297a13
Remove .vscode
rjurney Feb 17, 2025
2035d98
VERSION back to 0.8.4
rjurney Feb 17, 2025
f9f4bd7
Remove tutorials reference
rjurney Feb 17, 2025
9ddd6b2
VERSION is a Python thing, it belongs in python/
rjurney Feb 17, 2025
7065647
Include the README.md and LICENSE in the Python package
rjurney Feb 17, 2025
a6c7e91
Some classifiers for pyproject.toml
rjurney Feb 17, 2025
51e3e6d
Trying poetry install action instead of manual install
rjurney Feb 17, 2025
272be06
Removing SPARK_HOME
rjurney Feb 17, 2025
4587999
Returned SPARK_HOME settings
rjurney Feb 17, 2025
2422b22
Minimized the PR to just these files
rjurney Feb 17, 2025
073dced
Merge in rjurney/build-upgrades and in turn master
rjurney Feb 17, 2025
0a1faba
Created tutorials dependency group to minimize main bloat
rjurney Feb 17, 2025
c0d6d7b
Make motif.py execute in whole again
rjurney Feb 17, 2025
5bb4c26
Minor isort format and cleanup of download.py
rjurney Feb 17, 2025
99e6a4d
Minor isort format and cleanup of utils.py
rjurney Feb 17, 2025
662e197
Removed case sensitivity from the script - that was confusing people …
rjurney Feb 17, 2025
beaa35d
motif.py now matches tutorial code, runs and handles case insensitivity.
rjurney Feb 17, 2025
1bf4a9e
Regenerate poetry.lock
rjurney Feb 21, 2025
ef19784
Setup a 'graphframes stackexchange' comand.
rjurney Feb 21, 2025
4400cb4
Make graphframes.tutorials.motif use a checkpoint dir unique, and fro…
rjurney Feb 21, 2025
d549c56
Use spark.sparkContext.setCheckpointDir directly instead of instantia…
rjurney Feb 21, 2025
b970636
Using 'from __future__ import annotations' intsead of List and Tuple
rjurney Feb 21, 2025
3788941
Now retry three times if we can't connect for any reason in 'graphfra…
rjurney Feb 21, 2025
e95bbbe
Merge master
rjurney Feb 25, 2025
413a915
Merge branch 'master' of github.com:graphframes/graphframes
rjurney Mar 8, 2025
37ff13a
Add missing image
rjurney Mar 10, 2025
ae3c90a
Final docs fixes pre-release
rjurney Mar 10, 2025
00d1bfb
Minor newline fix
rjurney Mar 10, 2025
b3e2ce9
Fix bash/python messup
rjurney Mar 10, 2025
67e9830
Another newline
rjurney Mar 10, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
Removed case sensitivity from the script - that was confusing people …
…who just pasted or tried to run the code without a new SparkSession.
  • Loading branch information
rjurney committed Feb 17, 2025
commit 662e197960a424c1f58c151b663c46d9d63da6be
48 changes: 32 additions & 16 deletions python/graphframes/tutorials/stackexchange.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
# Build a Graph out of the Stack Exchange Data Dump XML files
"""Build a Graph out of the Stack Exchange Data Dump XML files."""

#
# Interactive Usage: pyspark --packages com.databricks:spark-xml_2.12:0.18.0
Expand Down Expand Up @@ -47,11 +47,9 @@ def split_tags(tags: str) -> List[str]:
# Initialize a SparkSession with case sensitivity
#

spark: SparkSession = (
SparkSession.builder.appName("Stack Exchange Graph Builder")
# Lets the Id:(Stack Overflow int) and id:(GraphFrames UUID) coexist
.config("spark.sql.caseSensitive", True).getOrCreate()
)
spark: SparkSession = SparkSession.builder.appName("Stack Exchange Graph Builder").getOrCreate()
sc = spark.sparkContext
sc.setCheckpointDir("/tmp/graphframes-checkpoints")

print("Loading data for stats.meta.stackexchange.com ...")

Expand Down Expand Up @@ -296,12 +294,23 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
)
print(f"Total distinct nodes: {nodes_df.count():,}")

# Now add a unique ID field
# Now add a unique lowercase 'id' field - standard for GraphFrames - moving the original...
# Stack Exchange Id to StackId
nodes_df = nodes_df.withColumnRenamed("Id", "StackId").drop("Id")

# Update the column list...
if "Id" in all_column_names:
all_column_names.remove("Id")
all_column_names += ["StackId"]
all_column_names = sorted(all_column_names)

# Add the UUID 'id' field for GraphFrames. It will go in edges as 'src' and 'dst'
nodes_df = nodes_df.withColumn("id", F.expr("uuid()")).select("id", *all_column_names)

# Now create posts - combined questions and answers for things that can apply to them both
posts_df = questions_df.unionByName(answers_df).cache()


#
# Store the nodes to disk, reload and cache
#
Expand Down Expand Up @@ -361,12 +370,12 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]

src_vote_df: DataFrame = votes_df.select(
F.col("id").alias("src"),
F.col("Id").alias("VoteId"),
F.col("StackId").alias("VoteId"),
# Everything has all the fields - should build from base records but need UUIDs
F.col("PostId").alias("VotePostId"),
)
cast_for_edge_df: DataFrame = src_vote_df.join(
posts_df, on=src_vote_df.VotePostId == posts_df.Id, how="inner"
posts_df, on=src_vote_df.VotePostId == posts_df.StackId, how="inner"
).select(
# 'src' comes from the votes' 'id'
"src",
Expand All @@ -378,6 +387,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
print(f"Total CastFor edges: {cast_for_edge_df.count():,}")
print(f"Percentage of linked votes: {cast_for_edge_df.count() / votes_df.count():.2%}\n")


#
# Create a [User]--Asks-->[Question] edge
#
Expand All @@ -388,7 +398,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
F.lit("Asks").alias("relationship"),
)
user_asks_edges_df: DataFrame = questions_asked_df.join(
users_df, on=questions_asked_df.QuestionUserId == users_df.Id, how="inner"
users_df, on=questions_asked_df.QuestionUserId == users_df.StackId, how="inner"
).select(
# 'src' comes from the users' 'id'
F.col("id").alias("src"),
Expand All @@ -402,6 +412,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
f"Percentage of asked questions linked to users: {user_asks_edges_df.count() / questions_df.count():.2%}\n"
)


#
# Create a [User]--Posts-->[Answer] edge.
#
Expand All @@ -412,7 +423,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
F.lit("Posts").alias("relationship"),
)
user_answers_edges_df = user_answers_df.join(
users_df, on=user_answers_df.AnswerUserId == users_df.Id, how="inner"
users_df, on=user_answers_df.AnswerUserId == users_df.StackId, how="inner"
).select(
# 'src' comes from the users' 'id'
F.col("id").alias("src"),
Expand All @@ -426,17 +437,18 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
f"Percentage of answers linked to users: {user_answers_edges_df.count() / answers_df.count():.2%}\n"
)


#
# Create a [Answer]--Answers-->[Question] edge
#

src_answers_df: DataFrame = answers_df.select(
F.col("id").alias("src"),
F.col("Id").alias("AnswerId"),
F.col("StackId").alias("AnswerId"),
F.col("ParentId").alias("AnswerParentId"),
)
question_answers_edges_df: DataFrame = src_answers_df.join(
posts_df, on=src_answers_df.AnswerParentId == questions_df.Id, how="inner"
posts_df, on=src_answers_df.AnswerParentId == questions_df.StackId, how="inner"
).select(
# 'src' comes from the answers' 'id'
"src",
Expand All @@ -450,6 +462,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
f"Percentage of linked answers: {question_answers_edges_df.count() / answers_df.count():.2%}\n"
)


#
# Create a [Tag]--Tags-->[Post] edge... remember a Post is a Question or Answer
#
Expand All @@ -472,6 +485,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
print(f"Total Tags edges: {tags_edge_df.count():,}")
print(f"Percentage of linked tags: {tags_edge_df.count() / posts_df.count():.2%}\n")


#
# Create a [User]--Earns-->[Badge] edge
#
Expand All @@ -482,7 +496,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
F.lit("Earns").alias("relationship"),
)
earns_edges_df = earns_edges_df.join(
users_df, on=earns_edges_df.BadgeUserId == users_df.Id, how="inner"
users_df, on=earns_edges_df.BadgeUserId == users_df.StackId, how="inner"
).select(
# 'src' comes from the users' 'id'
F.col("id").alias("src"),
Expand All @@ -494,6 +508,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
print(f"Total Earns edges: {earns_edges_df.count():,}")
print(f"Percentage of earned badges: {earns_edges_df.count() / badges_df.count():.2%}\n")


#
# Create a [Post]--Links-->[Post] edge... remember a Post is a Question or Answer
# Also a [Post]--Duplicates-->[Post] edge... remember a Post is a Question or Answer
Expand All @@ -505,15 +520,15 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
"LinkType",
)
links_src_edge_df: DataFrame = trim_links_df.join(
posts_df.drop("LinkType"), on=trim_links_df.SrcPostId == posts_df.Id, how="inner"
posts_df.drop("LinkType"), on=trim_links_df.SrcPostId == posts_df.StackId, how="inner"
).select(
# 'dst' comes from the posts' 'id'
F.col("id").alias("src"),
"DstPostId",
"LinkType",
)
raw_links_edge_df = links_src_edge_df.join(
posts_df.drop("LinkType"), on=links_src_edge_df.DstPostId == posts_df.Id, how="inner"
posts_df.drop("LinkType"), on=links_src_edge_df.DstPostId == posts_df.StackId, how="inner"
).select(
"src",
# 'src' comes from the posts' 'id'
Expand Down Expand Up @@ -557,6 +572,7 @@ def add_missing_columns(df: DataFrame, all_cols: List[Tuple[str, T.StructField]]
"count", F.format_number(F.col("count"), 0)
).show()


# +------------+------+
# |relationship| count|
# +------------+------+
Expand Down