如何用多进程加速Pandas DataFrame情感分析并创建新列?
解决多进程加速推特情感分析的问题
你遇到的'DataFrame' object has no attribute 'p'错误是因为错误调用了多进程池对象——p是你创建的Pool实例,不能通过j.p.apply这种方式调用,DataFrame本身没有p这个属性。以下是正确的多进程实现方案:
核心修正点
- 多进程环境中,
sonar、nlp这类模型对象无法直接在主进程初始化后传递给子进程,需要在每个子进程内部单独初始化,避免序列化报错。 - 改用
Pool的map/imap方法批量处理Tweet文本,而非直接对DataFrame调用apply。
完整代码实现
import multiprocessing as mp import pandas as pd # 子进程初始化函数:每个进程启动时加载一次模型,避免重复加载浪费资源 def init_worker(): global sonar, nlp # 在这里完成模型初始化,根据你的实际依赖调整 from asari.api import Sonar from transformers import pipeline nlp = pipeline("sentiment-analysis") sonar = Sonar() # 情感分析函数:仅接收文本参数,返回所需的分析结果 def analyze_tweet(text): asari_1 = sonar.ping(text) hug_1 = nlp(text) return ( asari_1['top_class'], asari_1['classes'][0]['confidence'], asari_1['classes'][1]['confidence'], hug_1[0]["label"], hug_1[0]["score"] ) if __name__ == '__main__': # 假设j是包含Tweet列的目标DataFrame # 初始化多进程池,可指定进程数(建议用CPU核心数-1,避免占满资源) with mp.Pool(initializer=init_worker, processes=mp.cpu_count()-1) as pool: # 批量处理所有Tweet文本,结果顺序与输入一致 results = pool.map(analyze_tweet, j['Tweet'].tolist()) # 将并行处理的结果拆分到对应列 j[['asar','asar neg','asar pos','hugposneg','hugscore']] = pd.DataFrame(results, index=j.index)
关键说明
- 模型复用:通过
initializer=init_worker让每个子进程仅加载一次模型,大幅提升处理效率。 - 内存优化:如果200万条数据直接用
map内存压力过大,可改用pool.imap_unordered(结果无序,但内存占用更低),示例如下:results = [] # 传递索引+文本,后续用于对齐原DataFrame顺序 for idx, res in enumerate(pool.imap_unordered(lambda x: (x[0], analyze_tweet(x[1])), enumerate(j['Tweet']))): results.append(res) # 按原索引排序后赋值 results.sort(key=lambda x: x[0]) j[['asar','asar neg','asar pos','hugposneg','hugscore']] = pd.DataFrame([r[1] for r in results], index=j.index)
内容的提问来源于stack exchange,提问作者talltreeee
相关产品推荐
相关产品推荐

