使用awswrangler.read_parquet读取S3表时触发无法启动新线程错误
根因分析
你遇到的RuntimeError: can't start new thread报错,本质是AWS Wrangler默认读取S3上的Parquet文件时,会启动多线程并行拉取多个文件/单个文件的多个分片,当并行任务量超过操作系统允许当前进程创建的最大线程数上限时就会触发报错。等待30分钟后恢复是因为之前占用的线程资源超时自动释放,线程数回到阈值以下即可正常运行。
解决方案
- 限制AWS Wrangler的并行线程数:在
wr.s3.read_parquet参数中添加use_threads配置,可直接设置为False禁用多线程,也可设置为较小的正整数(如2、4)主动控制最大线程量,修改后代码示例:aws_df = wr.s3.read_parquet(path=self._filepath, use_threads=2, **self._load_args) - 调高运行环境的线程上限:你当前使用EC2实例运行任务,可先执行
ulimit -u查看当前用户允许的最大进程/线程数,若数值较低可临时执行ulimit -u 4096调大阈值,需要永久生效则修改/etc/security/limits.conf文件,添加如下配置后重启会话即可:* soft nproc 4096 * hard nproc 8192 - 优化读取逻辑减少并行任务量:如果读取目标是大量小Parquet文件,可先将小文件合并为少量大文件再读取;如果是单个超大文件,可添加
s3_block_size参数调大S3分片大小(如设置为128 * 1024 * 1024即128MB),减少分片对应的并行线程需求。 - 排查进程内其他线程占用:检查你的项目代码中是否存在其他未正确释放的线程池、循环创建线程未回收的逻辑,修复这类问题可减少不必要的线程资源占用。
完整报错信息
aws_df = wr.s3.read_parquet(path=self._filepath, **self._load_args) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read_parquet.py", line 721, in read_parquet read_func=_read_parquet, paths=paths, version_ids=versions, use_threads=use_threads, kwargs=args File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read.py", line 145, in _read_dfs_from_multiple_paths return list(df for df in executor.map(partial_read_func, paths, versions)) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read.py", line 145, in <genexpr> return list(df for df in executor.map(partial_read_func, paths, versions)) File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/_base.py", line 586, in result_iterator yield fs.pop().result() File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/_base.py", line 432, in result return self.__get_result() File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/_base.py", line 384, in __get_result raise self._exception File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/thread.py", line 57, in run result = self.fn(*self.args, **self.kwargs) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read_parquet.py", line 495, in _read_parquet version_id=version_id, File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read_parquet.py", line 440, in _read_parquet_file source=f, read_dictionary=categories File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_read_parquet.py", line 40, in _pyarrow_parquet_file_wrapper return pyarrow.parquet.ParquetFile(source=source, read_dictionary=read_dictionary) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/pyarrow/parquet.py", line 201, in __init__ read_dictionary=read_dictionary, metadata=metadata) File "pyarrow/_parquet.pyx", line 1021, in pyarrow._parquet.ParquetReader.open File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_fs.py", line 569, in read self._fetch(self._loc, self._loc + length) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_fs.py", line 376, in _fetch self._cache = self._fetch_range_proxy(self._start, self._end) File "/home/ec2-user/anaconda3/lib/python3.7/site-packages/awswrangler/s3/_fs.py", line 359, in _fetch_range_proxy itertools.repeat(self._version_id), File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/_base.py", line 575, in map fs = [self.submit(fn, *args) for args in zip(*iterables)] File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/_base.py", line 575, in <listcomp> fs = [self.submit(fn, *args) for args in zip(*iterables)] File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/thread.py", line 160, in submit self._adjust_thread_count() File "/home/ec2-user/anaconda3/lib/python3.7/concurrent/futures/thread.py", line 181, in _adjust_thread_count t.start() File "/home/ec2-user/anaconda3/lib/python3.7/threading.py", line 847, in start _start_new_thread(self._bootstrap, ()) RuntimeError: can't start new thread
内容的提问来源于stack exchange,提问作者Miss.Saturn
相关产品推荐
相关产品推荐

