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

Sun Grid Engine集群中Python读磁盘出现IOError 110连接超时问题

IOError(110, 'Connection timed out') when reading files on Sun Grid Engine cluster

我在Sun Grid Engine(SGE)超级计算集群上运行Python脚本,脚本读取文件ID列表后分发到工作进程处理,每个输入文件会在磁盘写入输出。但现在工作函数minhash_text里出现了IOError(110, 'Connection timed out')错误——之前只有在网络请求严重延迟时才遇到过这个错误,但这次我只是读取磁盘文件。

想请教:读取磁盘为什么会出现连接超时错误?该怎么解决?

完整脚本(错误出现在minhash_text函数中)

from datasketch import MinHash
from multiprocessing import Pool
from collections import defaultdict
from nltk import ngrams
import json
import sys
import codecs
import config

cores = 24
window_len = 12
step = 4
worker_files = 50
permutations = 256
hashband_len = 4

def minhash_text(args):
    '''Return a list of hashband strings for an input doc'''
    try:
        file_id, path = args
        with codecs.open(path, 'r', 'utf8') as f:
            f = f.read()
        all_hashbands = []
        for window_idx, window in enumerate(ngrams(f.split(), window_len)):
            window_hashbands = []
            if window_idx % step != 0:
                continue
            minhash = MinHash(num_perm=permutations, seed=1)
            for ngram in set(ngrams(' '.join(window), 3)):
                minhash.update( ''.join(ngram).encode('utf8') )
            hashband_vals = []
            for i in minhash.hashvalues:
                hashband_vals.append(i)
                if len(hashband_vals) == hashband_len:
                    window_hashbands.append( '.'.join([str(j) for j in hashband_vals]) )
                    hashband_vals = []
            all_hashbands.append(window_hashbands)
        return {'file_id': file_id, 'hashbands': all_hashbands}
    except Exception as exc:
        print(' ! error occurred while processing', file_id, exc)
        return {'file_id': file_id, 'hashbands': []}

if __name__ == '__main__':
    file_ids = json.load(open('file_ids.json'))
    file_id_path_tuples = [(file_id, path) for file_id, path in file_ids.items()]
    worker_id = int(sys.argv[1])
    worker_ids = list(ngrams(file_id_path_tuples, worker_files))[worker_id]
    hashband_to_ids = defaultdict(list)
    pool = Pool(cores)
    for idx, result in enumerate(pool.imap(minhash_text, worker_ids)):
        print(' * processed', idx, 'results')
        file_id = result['file_id']
        hashbands = result['hashbands']
        for window_idx, window_hashbands in enumerate(hashbands):
            for hashband in window_hashbands:
                hashband_to_ids[hashband].append(file_id + '.' + str(window_idx))
    with open(config.out_dir + 'minhashes-' + str(worker_id) + '.json', 'w') as out:
        json.dump(dict(hashband_to_ids), out)

回答

这个错误看起来有点反直觉,但在SGE集群环境下,磁盘读取抛出连接超时通常和网络挂载的文件系统脱不了干系,而不是你以为的本地磁盘。大部分集群的存储都是NAS、NFS或者其他网络共享存储,你访问的“磁盘路径”本质上是通过网络连接到远程存储服务器的。当存储服务器负载过高、网络链路拥堵,或者计算节点和存储之间的连接不稳定时,就会触发这类超时错误。

另外,你的脚本用了24个并发进程同时读取文件,这种瞬间的高IO请求很容易压垮共享存储的带宽,进一步加剧连接超时的概率。

下面是几个针对性的解决办法:

  • 降低并发进程数:把cores变量从24调低,比如先试试8或12,减少同时访问存储的进程数量,给存储系统喘口气的机会。这是最快速的测试方法,能立刻验证是不是并发过高导致的问题。

  • 添加读取重试机制:在文件读取环节加上重试逻辑,遇到IO超时错误时等待几秒再重试,而不是直接放弃。修改minhash_text里的文件读取部分:

    import time
    
    def minhash_text(args):
        '''Return a list of hashband strings for an input doc'''
        try:
            file_id, path = args
            max_retries = 3
            retry_delay = 5  # 每次重试间隔5秒
            content = None
            for attempt in range(max_retries):
                try:
                    with codecs.open(path, 'r', 'utf8') as f:
                        content = f.read()
                    break
                except IOError as e:
                    if e.errno == 110 and attempt < max_retries - 1:
                        print(f" ! Retrying {file_id} (attempt {attempt+1}/{max_retries})...")
                        time.sleep(retry_delay)
                    else:
                        raise
            # 后续用content变量继续处理
            all_hashbands = []
            # ... 剩下的代码保持不变
    
  • 利用本地临时磁盘缓存文件:如果集群计算节点有本地临时存储(比如/tmp或者专门的SSD缓存盘),可以先把要处理的文件复制到本地,处理完成后再把输出写回共享存储。这样能彻底绕开网络存储的IO瓶颈:

    import shutil
    import tempfile
    import os
    
    def minhash_text(args):
        '''Return a list of hashband strings for an input doc'''
        try:
            file_id, path = args
            # 复制文件到本地临时路径
            with tempfile.NamedTemporaryFile(delete=False, encoding='utf8') as tmp_f:
                shutil.copyfileobj(codecs.open(path, 'r', 'utf8'), tmp_f)
                local_path = tmp_f.name
            # 读取本地文件
            with codecs.open(local_path, 'r', 'utf8') as f:
                content = f.read()
            # 处理完成后删除临时文件
            os.unlink(local_path)
            # ... 剩下的代码保持不变
    

    注意提前确认本地临时磁盘的空间是否足够放下要处理的文件。

  • 联系集群管理员排查底层问题:如果上面的方法都没效果,大概率是存储系统本身的问题——比如存储服务器故障、网络链路中断或者权限配置异常。联系管理员确认你访问的存储路径是不是网络共享存储,让他们帮忙检查存储节点的负载和健康状态。

  • 调整SGE作业参数:提交作业时可以指定更高的IO优先级,或者请求靠近存储节点的计算节点(如果集群支持的话)。比如用qsub的-l io=high参数(具体要看你的集群配置),或者通过-l node_type=storage_near这类参数调度到离存储更近的节点,减少网络传输延迟。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:30:34