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

如何在Databricks中通过并发作业实现并行处理加速文本操作

解决Databricks中Spark处理大规模文本数据的性能与序列化问题

问题背景

  • 需处理100万条文本记录,操作涵盖文本拼接、转换为spaCy Doc、提取动词等基础文本操作,未使用机器学习模型
  • 单条记录串行处理耗时1秒,总耗时约10天,需借助Databricks并发能力提速
  • 已尝试将代码从Pandas改为PySpark.Pandas,但仅提升30%性能(疑似因处理器升级),处理大规模数据时出现PicklingError: 无法序列化对象
  • 已获得分作业处理的推荐方向:创建带结果列和状态列的表、拆分10个作业处理未完成记录、使用Spark DataFrame,但不知如何落地

核心问题拆解

  1. 基于Spark DataFrame实现分批次作业+状态管理的并行处理
  2. 解决大规模数据处理时的Pickling序列化错误
  3. 优化spaCy在Spark分布式环境中的执行效率

一、解决PicklingError序列化问题

该错误源于spaCy的nlp对象无法被Spark默认的Pickle序列化机制处理,尤其是在分布式环境中传递对象时,可通过以下方式解决:

1. 用广播变量分发spaCy模型

在Driver端加载模型后,通过广播变量分发到所有Executor,避免重复序列化与加载:

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

# 初始化Spark会话
spark = SparkSession.builder.appName("SpaCyTextProcessing").getOrCreate()

# Driver端加载spaCy模型
nlp = spacy.load("en_core_web_sm")
# 广播模型到所有Executor
broadcast_nlp = spark.sparkContext.broadcast(nlp)

# 定义处理文本的UDF,使用广播的模型
def extract_verbs(text):
    doc = broadcast_nlp.value(text)
    return [token.lemma_ for token in doc if token.pos_ == "VERB"]

# 注册UDF
extract_verbs_udf = udf(extract_verbs, ArrayType(StringType()))

# 应用UDF到Spark DataFrame
processed_df = raw_df.withColumn("verbs", extract_verbs_udf(col("text_column")))

2. 避免UDF引用不可序列化对象

确保自定义UDF内仅使用广播变量、基础数据类型或可序列化对象,不要直接引用Driver端的非序列化对象。


二、实现分作业+状态管理的并行处理

基于Databricks默认支持的Delta Lake实现状态跟踪与分批次处理,支持断点续传:

1. 创建带状态的Delta表

将输入数据转为Delta表,并添加结果列与状态列:

-- 创建Delta表(已有表可跳过创建,直接ALTER添加列)
CREATE TABLE IF NOT EXISTS text_records (
    col1 STRING,
    col2 STRING,
    col3 STRING,
    col4 STRING,
    col5 STRING,
    result ARRAY<STRING>,
    status STRING
) USING DELTA LOCATION '/path/to/delta/table';

-- 初始化状态为'NaN',结果列为空
INSERT INTO text_records
SELECT col1, col2, col3, col4, col5, NULL AS result, 'NaN' AS status
FROM raw_input_data;

2. 分批次处理未完成记录

每次筛选status='NaN'的记录,按10万条/批次处理,同时更新状态:

from pyspark.sql.functions import lit

# 批次大小按100万/10设置
batch_size = 100000

# 获取未处理记录
unprocessed_df = spark.read.format("delta").load("/path/to/delta/table")\
    .filter(col("status") == "NaN").limit(batch_size)

# 更新状态为'InProgress'
unprocessed_df.withColumn("status", lit("InProgress"))\
    .write.format("delta").mode("overwrite").option("mergeSchema", "true")\
    .save("/path/to/delta/table")

# 执行文本处理
processed_batch_df = unprocessed_df.withColumn("result", extract_verbs_udf(col("text_column")))

# 更新状态为'Completed'并写入结果
processed_batch_df.withColumn("status", lit("Completed"))\
    .write.format("delta").mode("overwrite").option("mergeSchema", "true")\
    .save("/path/to/delta/table")

3. 并行调度作业

在Databricks Jobs中创建10个独立作业,每个作业执行上述批次逻辑;或直接利用Spark分布式计算能力,让集群根据资源自动分区并行处理(比手动拆分作业更高效)。


三、优化Spark处理性能

  1. 调整集群配置:增加Executor数量与内存,比如设置--num-executors 20、--executor-memory 8g,确保每个Executor有足够资源加载spaCy模型
  2. 分区优化:对Delta表按status字段分区,让Spark自动将数据分发到多个Executor并行处理
  3. 合并小文件:处理完成后合并Delta表小文件,提升后续查询效率:
OPTIMIZE text_records ZORDER BY status;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 04:15:40