80万句G2P音素生成与评分代码并行优化需求
并行化改造方案
针对你的80万条语句音素生成与评分场景,原始串行代码的瓶颈主要在音素生成的串行处理和共享triphones的低效遍历修改,以下是具体的并行化改造方案,同时修正代码中的逻辑错误:
先修正原始代码的几个问题
scoreSentence函数中的phonemes参数未被使用,属于冗余参数,已移除- 批量处理时错误地给剩余数据赋值,改为直接处理当前批次并保存
- 三重循环生成三音素效率低下,用
itertools.product替代
改造后完整代码
from app import getPhonemes import pandas as pd import sys from itertools import product from concurrent.futures import ProcessPoolExecutor, as_completed from multiprocessing import Manager, Lock def phonemize(sentence): tokens = sentence.split(' ') phonemes = getPhonemes(tokens) return '$'.join(phonemes) def generateTriphones(phonemes): # 用itertools.product替代三重循环,效率更高 return [' '.join(trip) for trip in product(phonemes, repeat=3)] def scoreSentence(phones, shared_triphones, lock): flag = 0 score = 0 tokens = phones.split('$') uniqueTokens = set(tokens) triphoneticTokens = [token for token in uniqueTokens if token.count(' ') > 1] # 加锁保证对共享triphones的操作原子性 with lock: for token in triphoneticTokens: # 遍历副本避免修改原列表导致的遍历异常 for triphone in list(shared_triphones): if token.find(triphone) != -1: score += 1 shared_triphones.remove(triphone) if not shared_triphones: flag = -1 break if flag == -1: break return score, flag def Process(fil): # 读取音素词典 with open('itudict/vocab.phoneme', 'r', encoding='utf-8') as file: data = [line.strip() for line in file] phonemes = data[4:] # 用Manager创建共享triphones和锁,支持多进程访问 manager = Manager() shared_triphones = manager.list(generateTriphones(phonemes)) lock = manager.Lock() # 读取原始数据 data = pd.read_csv(fil + '.csv') data = data.drop(['score', 'covered_vocab'], axis=1, errors='ignore') batch_size = 10000 total_batches = (len(data) + batch_size - 1) // batch_size print(f"Total batches to process: {total_batches}") for batch_idx in range(total_batches): print(f'Processing Batch: {batch_idx + 1}/{total_batches}') # 切分当前批次数据 start = batch_idx * batch_size end = min(start + batch_size, len(data)) sentence_batch = data.iloc[start:end] sentences = sentence_batch['sentence'].tolist() # 并行处理音素生成 phonemes_list = [] # 根据CPU核心数设置进程数,避免过载 with ProcessPoolExecutor(max_workers=4) as executor: future_to_sentence = {executor.submit(phonemize, sent): sent for sent in sentences} for future in as_completed(future_to_sentence): phonemes_list.append(future.result()) # 处理评分(由于triphones是共享消耗的,这里暂时串行,若要并行需更精细的锁控制) scores = [] flag = 0 for phones in phonemes_list: if flag == -1: scores.append(0) continue score, current_flag = scoreSentence(phones, shared_triphones, lock) scores.append(score) if current_flag == -1: flag = -1 print("Triphones exhausted, stopping processing") break # 给当前批次添加结果列并保存 sentence_batch['Phonemes'] = phonemes_list[:len(sentence_batch)] sentence_batch['score'] = scores[:len(sentence_batch)] sentence_batch.to_csv(f'{fil}_phonemized_{batch_idx + 1}.csv', index=False) if flag == -1: print("Early exit due to triphones exhaustion") break if __name__ == '__main__': if len(sys.argv) != 2: print("Usage: python script.py <input_file_prefix>") sys.exit(1) Process(sys.argv[1])
关键优化点说明
- 并行音素生成:使用
ProcessPoolExecutor并行处理语句的音素生成,充分利用多核CPU资源,这部分是原始代码的主要耗时点,并行后能大幅缩短时间。 - 共享triphones的线程安全:通过
multiprocessing.Manager创建共享列表和锁,保证多进程下对triphones的修改不会出现竞争条件,避免数据错乱。 - 三音素生成优化:用
itertools.product替代三重嵌套循环,代码更简洁且效率更高。 - 批量处理逻辑修正:改为按索引切分批次,直接处理当前批次并保存,避免原始代码中的赋值错误。
- 提前终止逻辑:当triphones耗尽时,直接终止后续处理,避免无用计算。
额外建议
- 若
getPhonemes本身支持批量处理,可进一步优化为批量传入tokens,减少进程间通信开销。 - 可根据你的CPU核心数调整
ProcessPoolExecutor的max_workers参数(通常设为CPU核心数或核心数+1)。 - 若评分部分仍有性能瓶颈,可考虑将triphones转换为集合或前缀树(Trie),优化匹配速度,因为当前的遍历匹配是O(n)复杂度,数据量262万时非常耗时。
内容的提问来源于stack exchange,提问作者Umair Afzal
相关产品推荐
相关产品推荐

