基于Spark与JSON实现文本文件键匹配及恐怖类长书评书籍筛选
Hey Jean, let’s walk through solving this Spark problem step by step. I’ll cover both the core filtering task you need and how to handle JSON key lookups/comparisons along the way.
First, let’s make sure we’re loading your data correctly. From your description, BookInformation.data has entries with userName, price, categories, and title—either as line-delimited JSON or CSV. I’ll cover both common scenarios:
If it’s line-delimited JSON (each row is a full JSON object like your example):
from pyspark.sql import SparkSession from pyspark.sql.functions import col, size, split, lower, get_json_object spark = SparkSession.builder.appName("HorrorBookFilter").getOrCreate() # Load the book info directly as JSON book_info_df = spark.read.json("path/to/BookInformation.data")If it’s CSV (comma-separated with headers):
# Load with inferred schema (adjust if price isn't auto-detected as integer) book_info_df = spark.read.csv("path/to/BookInformation.data", header=True, inferSchema=True)
For your review file (let’s call it BookReviews.data), I’ll assume it has at least two columns: title (to match the book info) and review_text (the actual review content). Load it similarly:
reviews_df = spark.read.csv("path/to/BookReviews.data", header=True, inferSchema=True)
To find reviews with more than 100 words, we’ll split the review text by whitespace and count the elements:
# Split review text into words, count them, and filter for >100 filtered_reviews_df = reviews_df.filter( size(split(col("review_text"), "\\s+")) > 100 )
Note: If your reviews use non-whitespace word separators, adjust the regex in split() to match (e.g., [\\s,.!]+ for punctuation).
Now we’ll join the filtered reviews with the book information to get category data, then keep only horror books:
# Join on title (use a unique ID if titles aren't unique!) result_df = filtered_reviews_df.join(book_info_df, on="title", how="inner") \ # Make the filter case-insensitive (covers "Horror", "horror", "HORROR") .filter(lower(col("categories")).contains("horror")) \ # Select only the columns you care about .select("title", "userName", "price", "review_text") # View the result or save it result_df.show(truncate=False) # result_df.write.csv("path/to/filtered_horror_books.csv", header=True)
If you’re working with raw JSON strings (instead of parsed DataFrames), here’s how to access and compare keys:
Extract values from JSON strings: Use
get_json_objectto pull specific keys:# If you loaded book info as raw text (each line is a JSON string) book_info_text_df = spark.read.text("path/to/BookInformation.data") parsed_book_df = book_info_text_df \ .withColumn("userName", get_json_object(col("value"), "$.userName")) \ .withColumn("categories", get_json_object(col("value"), "$.categories")) \ .withColumn("price", get_json_object(col("value"), "$.price").cast("integer")) \ .drop("value")Compare key values: Once parsed, you can filter or compare just like any other column:
# Example: Filter books priced over $20 pricey_horror_df = result_df.filter(col("price") > 20) # Example: Check if a key exists (using `isNotNull`) valid_books_df = book_info_df.filter(col("userName").isNotNull())
Let me know if you need tweaks based on your actual file formats—sometimes CSV/JSON has quirks like extra spaces or nested categories. Happy coding!
内容的提问来源于stack exchange,提问作者JeanP

