You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame每行句子词干提取与分句预处理需求

Hey there! Let's walk through both of your Spark DataFrame text processing tasks step by step. I'll use PySpark for the examples (since it's super accessible), but you can adapt the logic to Scala if that's your stack.

需求1:对每行的每个句子执行词干提取操作

First, we need to split each row's text into individual sentences, then apply stemming to every word in each sentence. We'll use NLTK for sentence splitting and stemming (just make sure you install NLTK on all your Spark executors first!). Also, we'll use a UDF to handle per-row processing, and broadcast the stemmer to avoid reinitializing it on every executor.

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, col
from pyspark.sql.types import StringType
import nltk
from nltk.tokenize import sent_tokenize
from nltk.stem import PorterStemmer

# Initialize Spark session
spark = SparkSession.builder.appName("TextProcessing").getOrCreate()

# Download NLTK resources (run once on each executor)
nltk.download('punkt')

# Broadcast the stemmer to all executors
stemmer = PorterStemmer()
broadcast_stemmer = spark.sparkContext.broadcast(stemmer)

# Define UDF for processing each row
def stem_sentences(text):
    # Split text into sentences
    sentences = sent_tokenize(text)
    stemmed_sentences = []
    for sent in sentences:
        # Split sentence into words
        words = sent.split()
        # Stem each word
        stemmed_words = [broadcast_stemmer.value.stem(word) for word in words]
        # Join back into a sentence
        stemmed_sent = ' '.join(stemmed_words)
        stemmed_sentences.append(stemmed_sent)
    # Join all stemmed sentences back (adjust separator as needed)
    return ' '.join(stemmed_sentences)

# Register UDF
stem_sentences_udf = udf(stem_sentences, StringType())

# Test DataFrame
df = spark.createDataFrame([
    ("This is a sample sentence. I love Spark processing!",),
    ("Stemming helps reduce word forms. Running runs ran are all run!",)
], ["text"])

# Apply the UDF
result_df = df.withColumn("stemmed_text", stem_sentences_udf(col("text")))

# Show results
result_df.show(truncate=False)

Quick Notes:

  • If you prefer Spark's MLlib over NLTK, you can chain RegexTokenizer → SnowballStemmer in a pipeline, but you'll need to handle sentence splitting first (MLlib's tokenizers target words, not sentences by default).
  • Ensure NLTK's punkt tokenizer is available on all workers—package it with your job or use a bootstrap script to install it.
需求2:每行句子的文本预处理(小写+去标点+格式化为[]包裹的句子集合)

This task requires splitting sentences, cleaning each one (lowercase + remove punctuation), then wrapping each cleaned sentence in [] and concatenating them. We'll use regex to handle splitting and punctuation removal efficiently.

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf, col
from pyspark.sql.types import StringType
import re

# Initialize Spark session
spark = SparkSession.builder.appName("TextPreprocessing").getOrCreate()

# Define UDF for the preprocessing
def preprocess_text(text):
    # Step 1: Split text into sentences (handles cases like "text.I" with no space after punctuation)
    sentences = re.split(r'(?<=[.!?])\s*', text)
    cleaned_sentences = []
    for sent in sentences:
        # Skip empty strings from split
        if not sent.strip():
            continue
        # Step 2: Convert to lowercase
        lower_sent = sent.lower()
        # Step 3: Remove all punctuation (keep letters, numbers, spaces)
        no_punc_sent = re.sub(r'[^\w\s]', '', lower_sent)
        # Step 4: Remove extra spaces (trim leading/trailing, collapse multiple spaces)
        cleaned_sent = re.sub(r'\s+', ' ', no_punc_sent).strip()
        # Step 5: Wrap in [] if the sentence isn't empty
        if cleaned_sent:
            cleaned_sentences.append(f"[{cleaned_sent}]")
    # Join all wrapped sentences
    return ''.join(cleaned_sentences)

# Register UDF
preprocess_text_udf = udf(preprocess_text, StringType())

# Test DataFrame with your example
df = spark.createDataFrame([
    ("This is a text.I want to split!",),
    ("Hello World! How are you? I'm doing great.",)
], ["text"])

# Apply the UDF
result_df = df.withColumn("processed_text", preprocess_text_udf(col("text")))

# Show results
result_df.show(truncate=False)

Example Output:

For your input This is a text.I want to split!, you'll get exactly what you requested:

[this is text][i want to split]

Quick Notes:

  • The regex (?<=[.!?])\s* uses a positive lookbehind to split sentences only after punctuation marks, which handles missing spaces between sentences correctly.
  • Adjust the punctuation removal regex [^\w\s] if you need to keep specific characters (e.g., hyphens or apostrophes).

内容的提问来源于stack exchange,提问作者user5801828

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:40:08