Python ThreadPoolExecutor или ProcessPoolExecutor не хватает памятиPython

Программы на Python
Anonymous
Python ThreadPoolExecutor или ProcessPoolExecutor не хватает памяти

Сообщение Anonymous »

Я написал интерфейс, который в основном используется для запроса данных и возврата их во внешний файл данных после обработки. В процессе запроса данных я использую несколько потоков, а обработка данных использует несколько процессов.
Мой сервис развернут на K8S, объем запрашиваемых пользователями данных иногда бывает большим (возможно, десятки от тысяч до миллионов), а объем памяти контейнера часто превышает 10 ГБ после нескольких запросов. 10G: я добавляю ограничения в контейнер, контейнер автоматически перезапускается, иногда появляется следующее сообщение об ошибке: Параллельное будущее. Процесс.
BrokenProcessPool: процесс в пуле процессов был внезапно завершен во время выполнения или ожидания будущего. Процесс в пуле процессов был внезапно завершен во время выполнения или ожидания Future.
Я попробовал gc.collect, разбиение на фрагменты и запись данных в файл, как показано на рисунке. в коде ниже. Как мне оптимизировать, чтобы решить эту проблему?
@download_bp.route('/dataDownload', methods=['POST'])
@siwa.doc(body=GetChartDataBody, tags=['数据下载'], summary='数据下载')
def download_file(body: GetChartDataBody):
file_path = get_chart_data2file(body)['data']

@after_this_request
def remove_file(response):
try:
os.remove(file_path)
except Exception as error:
logger.error("Error removing or closing downloaded file handle", error)
return response

return send_file(file_path, mimetype='application/csv')


import os
import gc
import time
from flask import request
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from functools import partial
from ..utils import get_login_user
from ..log import logger
from ..dtos.request_body import GetChartDataBody
from ..services.get_vars import conn_influxdb
import pandas as pd
from ..dtos.resp_result import RespResult
from ..exception import BizException
from ..settings import settings
from datetime import datetime, timedelta
from ..services import pd_data_nan2none
from . import user_file_path
import pickle
import uuid
import shutil

def get_chart_data2file(body: GetChartDataBody):
"""
多线程查询,多进程处理
:param body:
:return:
"""
gc.collect()
terminal_id = body.terminal_id
vars_list = body.vars_list
start_time_str = body.start_time
end_time_str = body.end_time
table_name = terminal_id + '_xcp'
record_time_start = time.time()
user_name = get_login_user(request)
user_directory = os.path.join(user_file_path, f"{user_name}_{str(uuid.uuid4())}")
os.makedirs(user_directory, exist_ok=True)
result_df = gen_result_df(user_directory, table_name, vars_list, start_time_str, end_time_str, True)
dtype_dict = {col: 'float32' for col in result_df.columns if col != 'timestamps'}
dtype_dict['timestamps'] = result_df['timestamps'].dtype
result_df = result_df.astype(dtype_dict)
file_path = os.path.join(user_file_path, settings.FILE_NAME)
result_df.to_csv(file_path, index=False, chunksize=10000)
record_time_end = time.time()
spend_time = record_time_end - record_time_start
logger.info(f'{user_name} download all finished: terminal_id:{terminal_id}, query_vars:{vars_list}, query_time_range:{start_time_str}至{end_time_str},spend_time: {spend_time}s')
del result_df
gc.collect()
return RespResult.success(file_path)

def gen_result_df(user_directory, table_name, vars_list, start_time_str, end_time_str, is_download=False):
"""
生成处理完的df
:param user_directory
:param table_name:
:param vars_list:
:param start_time_str:
:param end_time_str:
:param is_download: 降频or 不降
:return:
"""
logger.info(f'start query and handle data:{time.time()} ')
time_groups = gen_time_group(start_time_str, end_time_str)
with ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(partial(query_database, user_directory), vars_list, table_name, start_time, end_time) for start_time, end_time in time_groups]
results = [future.result() for future in futures]
del futures
gc.collect()

if not any(results):
shutil.rmtree(user_directory)
raise BizException(201, '该时间范围内无数据,请重新选择时间')

file_paths = [os.path.join(user_directory, file) for file in os.listdir(user_directory)]

with ProcessPoolExecutor(max_workers=2) as executor:
start = 0
cnt = len(results)
df_chunks = []
while start < cnt:
df_chunks.extend(executor.map(partial(handle_data, is_download=is_download), file_paths[start:start + 3]))
start += 3
gc.collect()

result_df = pd.concat(df_chunks, ignore_index=True)
result_df.sort_values(by='timestamps', ascending=True, inplace=True)
shutil.rmtree(user_directory)
del results, time_groups, df_chunks
gc.collect()

return result_df

def gen_time_group(start_time_str, end_time_str):
"""
将时间分块,按时间多线程分块查询
:param start_time_str:
:param end_time_str:
:return:
"""
start_time = datetime.strptime(start_time_str, '%Y-%m-%d %H:%M:%S')
end_time = datetime.strptime(end_time_str, '%Y-%m-%d %H:%M:%S')
day = timedelta(days=0.5)
time_ranges = []
current_time = start_time
while current_time < end_time:
if current_time + day

Подробнее здесь: https://stackoverflow.com/questions/787 ... -of-memory

Вернуться в «Python»