PySpark行删除与合并问题:解析文件后数据行错位需修正
Problem Overview
You're trying to parse a text file where URLs and their associated text lines sit on separate rows, but your current PySpark code is spitting out empty rows and misaligned column values. Here's your original code:
File = "somehtml.file" Data = spark.read.text(File) df_file = Data.select( regexp_extract("col1", '(.*?)', 0).alias("somedata"), regexp_extract("col1", '(.*?)', 0).alias("somedata2") )
And the messy output you're seeing:
+--------------------+--------------------+ | somedata | somedata2 | +--------------------+--------------------+ |http://sweersdsh.ru....| | | |helo my name lololol...| | | | |http://qweuiewjk.ru....| | | |helo my name alallal...| +--------------------+--------------------+
Your goal is to pair each URL with its matching text line, with no empty rows cluttering the result:
+--------------------+--------------------+ | somedata | somedata2 | +--------------------+--------------------+ |http://sweersdsh.ru....|helo my name lololol...| |http://qweuiewjk.ru....|helo my name alallal...| +--------------------+--------------------+
Root Cause
The core issue is that your text file stores URLs on one line and their corresponding text on the next—plus there are random empty lines thrown in. Your current code processes every line in isolation, so each row only has one non-empty value (either the URL or the text), and empty lines are kept in the dataset.
Solution
Let's fix this step by step, with code that cleans the data, groups URLs with their text, and removes empty rows:
Step 1: Import Required PySpark Functions
from pyspark.sql import functions as F from pyspark.sql.window import Window
Step 2: Read & Clean the Raw Data
First, we'll read the file, add line numbers to preserve order, and filter out empty lines right away:
# Read the file (note: the default column name for text files is "value", not "col1") df = spark.read.text("somehtml.file") # Add line numbers to maintain the original order of lines, then filter out empty rows df_clean = df.withColumn("line_num", F.monotonically_increasing_id()) \ .filter(F.col("value") != "")
Step 3: Identify URL Lines & Create Groups
We'll mark which lines are URLs, then use a window function to group each URL with the text line immediately after it:
# Tag rows that contain a URL (matches lines starting with http:// or https://) df_with_tags = df_clean.withColumn( "is_url", F.col("value").rlike("^http[s]?://") ) # Create a group ID: each URL starts a new group, text lines inherit the group from the last URL window_spec = Window.orderBy("line_num") df_with_groups = df_with_tags.withColumn( "group_id", F.sum(F.when(F.col("is_url"), 1).otherwise(0)).over(window_spec) )
Step 4: Aggregate to Pair URLs & Text
Finally, we'll group by the group_id and extract the URL and text into their respective columns:
# Pair each URL with its matching text line final_df = df_with_groups.groupBy("group_id") \ .agg( F.max(F.when(F.col("is_url"), F.col("value"))).alias("somedata"), F.max(F.when(~F.col("is_url"), F.col("value"))).alias("somedata2") ) \ .drop("group_id") # View the cleaned result final_df.show(truncate=False)
Explanation
- Filtering empty lines: We eliminate blank rows upfront to avoid clutter in the final output.
- Line numbers:
monotonically_increasing_id()ensures we don't lose the original order of lines—critical for pairing URLs with the correct text. - Group IDs: The window function sums up URL markers to create a unique group for each URL-text pair, so every text line is linked to the most recent URL above it.
- Aggregation: Using
max()withwhen()lets us pull the URL and text from each group into separate, aligned columns.
This code will produce exactly the clean, paired output you're looking for.
内容的提问来源于stack exchange,提问作者vamper1234

