В большинстве компаний продукт выкатывают на разных ландшафтах: пишут код на локальной машине разработчика, затем код разворачивают на тесте, после чего на предпродуктивной системе. Если на этих этапах не возникло критичных замечаний и приложение работает штатно, то оно разворачивается на продуктиве.

Но, как правило, ландшафты имеют разные вычислительные ресурсы. На продуктив их выделяют больше, а на тест и предпродуктив — существенно меньше. Причем тестировать на них бизнес‑пользователи обычно хотят большую часть функциональности.

Привет, Хабр! Меня зовут Александр, и я занимаюсь развитием системы расчета налога на дополнительный доход в ИТ‑кластере. В этой статье расскажу, как без боли для себя и для пользователя выгрузить большие таблицы из базы данных (БД) и предоставить их пользователю в xls‑формате. Рассмотрю вариант выгрузки всех необходимых данных в одном файле без кэширования и объясню, почему такое подход не будет работать.

 Ландшафты
Ландшафты

Наша тестовая система имеет 16 Gb оперативной памяти и 6 CPU, в то время как продуктивная система — 64 Gb RAM и 12 CPU.


Test

Stage

Prod

CPU

6

6

12

RAM (Gb)

16

16

64

Из‑за этого возникают ситуации подобно той, с которой столкнулись мы. При этом эти ресурсы делятся между фронтом и бэком.

Расскажем немного о нашем кейсе. Мы занимается поддержкой и развитием системы, которая считает налог на дополнительный доход (НДД). Основная часть функциональности — это расчет НДД, просмотр и скачивание декларации в формате xml/xlsx, расчетных регистров для сверки данных. У пользователей есть действие, которое они делают каждое закрытие и многократно — это выгрузка расчетных регистров в одном файле xls. В нем зафиксированы определенные форматы данных. За счет того, что продуктивная система имеет больше вычислительных ресурсов, выгрузка в продуктивной системе проходит штатно для пользователя. Но в тестовой системе мы столкнулись с тем, что Pod в Deckhouse перезагружается.

Проанализировали ситуацию и пришли к выводу, что перезагрузка происходит, когда доходит лимит по памяти: ее не хватает, поэтому Deckhouse перезагружает Pod. Стоит отметить, что выгрузка файла была достаточно долгая даже в продуктиве. Занимало это примерно 30 минут — столько мы держали http‑соединение, когда пользователь на фронте нажимал кнопку выгрузки и шел запрос на бэкенд. При этом в тестовой системе эта выгрузка почти всегда не проходила и пользователь получал ошибку после такого долгого ожидания.

Первоначальная версия выгрузки

Изначальная версия выглядела так:

 Схема первой версии выгрузки
Схема первой версии выгрузки

Недостаток такого формирования файла — это удержание большого объема данных в памяти. Изначально мы формировали список регистров, которые должны выгрузить. Это можно получить из модели данных. Далее создали пустой xls‑файл и держали его в памяти. Тут важно отметить, что мы будем писать данные в определенном формате, то есть в памяти мы будем хранить этот большой объем данных. Далее мы пробегаемся по списку регистров, запрашиваем регистр у БД и сохраняем его как новый лист в файле. Весь этот цикл работает, пока у нас есть регистры в списке. После того как регистры закончились, мы считаем, что файл сформирован, и отправляем его пользователю.

from io import BytesIO

from flask import Response
from pandas import ExcelWriter

from .api.app import crud
from .api.models.registers._base_model_registers import RegistersModel
from .api.views.base_api_resource import BaseApiResource
from .utils.excel_format import ExcelFormat
from .utils.report_excel_writer import ReportExcelWriter

class ExportAllRegisters(BaseApiResource):
	def get(self) -> Response:
		file = BytesIO()

		with ExcelWriter(file, engine="xlsxwriter") as writer:
			formats = ExcelFormat(writer.book)
			list_models = [item for item in RegistersModel.__subclasses__()]
			for register_model in list_models:
				register_name = register_model.__tablename__

				actual_df = crud.get_actual_df(
					table_name=register_model.__tablename__,
					bu=bu,
					version=version,
					year=year,
					quarter=quarter,
				)
				report_writer = ReportExcelWriter(
					crud=crud,
					df=actual_df,
					db_object_name=register_name,
					writer=writer,
					formats=formats
				)
				report_writer.write()

		file.seek(0)

		return Response(file, mimetype="application/octet-stream")

Приложение у нас написано на flask, и мы используем отдельные классы, чтобы реализовать endpoints.

  • RegistersModel — базовый класс для вычислительных регистров.

  • BaseApiResource — базовый класс для endpoints.

  • ExcelFormat — отдельный класс, который хранит форматы для xls‑файлов.

  • ReportExcelWriter — отдельный класс для реализации записи данных в xls, в качестве параметров он получает crud, чтобы можно было получить типы данных из БД и корректно их записать в xls.

    Также:

  • df — сам датафрейм с данными;

  • db_object_name — наименование регистра;

  • writer — объект ExcelWriter и необходимые форматы formats.

Далее файл возвращается фронту. Как видно, в такой выгрузке мы всегда храним в памяти форматированные данные xls, которые могут быть достаточно большими.

Мы используем класс RegisterModel как базовый класс для регистров. Для выгрузки xls‑файлов используются классы: ExcelFormat — хранит форматы необходимые для выгрузки xlsx, ReportExcelWriter — осуществляет запись с учетом типов данных.

В таком решении мы имеем одну большую проблему при выгрузке. Из‑за того что мы храним весь файл в памяти, памяти нам не хватает, Pod перезагружается каждый раз, когда мы пытаемся выгрузить файл. Помимо этого, процесс выгрузки висит 30–40 минут.

Оптимизация

Первым вариантом было разбить хранение регистров в отдельных файлах и упаковать их в архив. Далее архив сохранить на диск, чтобы каждый раз его не генерировать.

Какие библиотеки необходимо добавить:

from datetime import datetime, timezone
from os import remove
from typing import Generator
from zipfile import ZIP_LZMA, ZipFile

from flask import Response, stream_with_context

Далее мы создадим имя для файла.

curr_time = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")
zip_file = ROOT_FOLDER / "tmp" / f"all_registers_{curr_time}.zip"

Как вы помните, когда мы запрашиваем get_actual_df, мы передаем ему балансовые единицы, год, версию данных и квартал. Мы не стали добавлять их в название файла, чтобы хоть как‑то отличать файлы. Для начала нам нужно было понять, будет ли работать теория, при которой мы сохраняем файл на диск. При каждом обращении к endpoint этот файл удалялся.

Далее мы перенесли формирование файла в блок формирования архива.

with ZipFile(self.zip_file, mode="w", compression=ZIP_LZMA, compresslevel=9) as zf:
	zf.writestr(
		data=register_file.getvalue(),
		zinfo_or_arcname=f"{register_name}.xlsx",
		compress_type=ZIP_LZMA,
		compresslevel=9,
	)

Думая, что, упаковав данные в архив мы решим проблему, мы ошиблись. Основная проблема была в том, что мы по‑прежнему запрашивали из БД большие порции данных. Мы их никуда не кэшировали, а хранили в памяти, только в другом формате. И именно поэтому, когда мы запрашивали большие регистры, у нас снова были большие скачки по памяти.

Тогда мы пошли анализировать таблицы, которые нам необходимо выгрузить из БД и отдать пользователю, чтобы понять, как мы можем запрашивать у БД таблицы небольшими порциями. Данные у нас хранятся как по месяцам, так и по кварталам. Есть таблицы с месяцем, а есть таблицы, в которых вместо месяца квартал. В квартальных таблицах данных, как правило, меньше, а в месячных — намного больше. Можно запрашивать из БД не весь квартал для таких таблиц, а только месяц, а также разбить большой файл на три и сохранять каждый отдельный файл на диск.

Напишем небольшой метод, который будет возвращать нам доступные месяца для текущей таблицы — get_available_months. Он будет смотреть в модели — если есть колонка «месяц», то возвращаются месяца для текущего квартала, если нет — список будет пуст.

Чтобы оценить, сколько сейчас записей в таблице для текущего квартала, сделаем метод get_rows_count. Он будет возвращать количество строк для текущего квартала, которое сейчас есть в базе данных, для требуемой таблицы.

Получим, сколько всего строк на данный момент в таблице, чтобы сравнить с порогом.

rows_count = self._get_rows_count(model)

И проверим это.

if rows_count > self.single_file_threshold:
	months_for_split = self._get_available_months(model)

Далее мы каждый отдельный файл можем сохранить на диск и готовую папку упаковать в архив.

zs = ZipStream(compress_type=ZIP_DEFLATED, compress_level=9)
zs.add_path(xlsx_files_folder)
with open(zip_file_name, "wb") as zfile:
	zfile.writelines(zs)

Также все данные действия лучше поместить в try:

try:
	…
except Exception:
	raise
finally:
	shutil.rmtree(xlsx_files_folder)
	log_info(f"Удалена временная папка {xlsx_files_folder} для размещения XLSX.")
	alive_tables_helper.clear_gen_zip_alive_status()

При любом исходе нам нужно будет удалить временную папку и убрать статус из таблицы.

Важно понимать, что в системе расчета НДД главным является сам расчет НДД и выгрузка декларации. А выгрузка всех регистров — это дополнительный приятный функционал. Поэтому важно, что в процессе расчета выгрузка архива с регистрами недоступна. А если расчет еще ни разу не запускался для данного периода, то выгружать нечего. Помимо всего этого пользователю хотелось бы понимать, на каком этапе находится выгрузка архива, а не просто 30 минут смотреть на крутящееся колесико.

Начнем с того, что создадим статусы выгрузки. Это позволит нам оповещать пользователя о текущих шагах.

from enum import Enum

class GeneratedZipSteps(Enum):
	calculation_alive = "Выполняется пайплайн расчёта декларации с текущими параметрами. Скачивание zip недоступно."
	start_creating_zip = "Выполнен запуск формирования zip."
	download_register_from_db = "Регистр скачивается с БД."
	save_register_on_disk = "Регистр сохраняется на диск."
	create_zip = "Создается архив с регистрами."
	unable_to_generate_file = "Невозможно сформировать zip из-за нехватки дискового пространства."
	pipeline_has_not_been_run_yet = "Пайплайн расчета периода еще ни разу не был запущен."
	pipeline_has_been_run_while_zip_generating = "В процессе формирования zip был запущен пайплайн расчета"

У нас уже была таблица, которая хранила статусы расчета. А таблицы, которая бы хранила статус выгрузки архива, не было. Давайте ее тоже заведем. Так мы сможем получать статус из базы данных, когда другой пользователь захочет сделать такую же выгрузку.

class User(Base):
	__tablename__ = "gen_zip_alive_status"

	code_bu = Column(String(10), primary_key=True)
	version = Column(Integer,primary_key=True)
	year = Column(Integer,primary_key=True)
	quarter = Column(Integer, primary_key=True)
	status = Column(String(100), nullable=False)
	register = Column(String(20), nullable=False)

Так как наши основные параметры — это БЕ, версия данных, год и квартал, то их и добавим в составной ключ. Также будем хранить текущий статус и регистр.

Теперь создадим класс, который поможет проверять, запущен ли текущий расчет или формирование архива.

class AliveTablesHelper:
	def __init__(self, bu: str, version: int, year: int, quarter: int):
		self.bu = bu
		self.version = version
		self.year = year
		self.quarter = quarter

Добавим в него методы:

  • check_calc_alive — возвращает True/False в зависимости от того, запущен расчет с необходимыми параметрами или нет;

  • check_gen_zip_alive — аналогично проверяет запущено ли формирование zip;

  • get_current_status — возвращает текущий статус и регистр для формирования zip‑архива;

  • clear_gen_zip_alive_status — выполняет очистку статуса для текущих параметров формирования zip;

  • write_gen_zip_alive_status_to_db — пишет текущий статус выгрузки для заданных БЕ, года, версии, квартала.

Класс больше работает с таблицей выгрузки архива, чем с таблицей статусов расчета, так как от таблицы статуса расчета нам нужен только статус, а с таблицей выгрузки архива нам необходимо будет работать в новой версии endpoint.

Теперь определим, что должен уметь делать этот класс.

def _get_file(self) -> Union[Path, None]:

	mask = self.file_mask + "_*" + self.file_extension
	list_of_files = glob(
	os.path.join(
		self.root_path,
		mask,
	)	
	)
	# файл всегда должен быть один
	sorted_files = sorted(list_of_files, key=os.path.getmtime, reverse=True)
	if sorted_files:
		return Path(sorted_files[0])

	return None

Во‑первых, нужно искать необходимый файл по маске на диске. Если файл уже есть — возможно, нам не потребуется его выгружать снова. Выгрузка зависит от того, был ли перезапущен расчет после последней выгрузки файла и сохранения его на диск. Если расчет перезапускался — файл перестал быть актуальным. Именно поэтому в название файла мы добавим метку времени.

Теперь напишем небольшой метод проверки существования файла.

def is_file_exist(self) -> bool:

	if self._get_file():
		return True
	return False

Если файл существует и он актуальный — выгрузка для пользователя ускорится в разы, так как на диске будет необходимый файл и мы сразу можем вернуть его пользователю.

def get_file_name_size(self) -> Tuple[str, int]:

	file_path = self._get_file()
	if not file_path:
		raise FileNotFoundError("Файл не существует.")
	file_size = os.path.getsize(file_path)
	file_name = file_path.name

	return file_name, file_size

Также мы должны получать файл и его размер, чтобы при необходимости отдавать эту информацию фронту.

Помимо того что мы храним текущий статус расчета в отдельной таблице, мы храним и статус расчета по каждому отдельному регистру. И его мы тоже должны уметь получать. Поэтому напишем небольшой метод, который будет получать время последнего шага расчета из этой таблицы — get_pipeline_last_time_run. Он будет ходить в эту таблицу и брать максимальное время расчета последней таблицы в пайплайне запуска расчета НДД.

Это поможет выяснить, является ли тот архив, который сейчас лежит на диске, актуальным.

def is_zip_actual(self) -> bool:

	# дату файла возьмем из названия файла
	file_path = self._get_file()
	if not file_path:
    	return False
	file_date_str = file_path.name.replace(self.file_mask + "_", "").replace(self.file_extension, "")
	file_date = datetime.strptime(file_date_str, self.date_format)

	pipeline_status_last_run = self.get_pipeline_last_time_run()

	if not pipeline_status_last_run:
		return False

	if file_date > pipeline_status_last_run:
		return True

	return False

Теперь мы можем написать основную функцию этого класса — это генерация архива.

def _generate_zip(self, app_context: "AppContext") -> None:

Метод будет принимать текущий контекст приложения, чтобы мы могли запустить его в потоке. Теперь разберем работу метода подробнее.

curr_time = datetime.now(timezone.utc)
curr_time_formatted = curr_time.strftime(self.date_format)

xlsx_files_folder = self.root_path / Path(f"{self.file_mask}_{curr_time_formatted}")
zip_file_name = self.root_path / Path(f"{self.file_mask}_{curr_time_formatted}" + self.file_extension)

app_context.push()
file_path = self._get_file()
if file_path and os.path.isfile(file_path):
	os.remove(file_path)
	log_info(f"Файл {file_path} удален.")

xlsx_files_folder.mkdir(parents=True, exist_ok=True)
log_info(f"Создана временная папка {xlsx_files_folder} для размещения XLSX.")

Нам нужно определить текущее время, которое мы добавим к названию файла. Сначала мы создадим на диске папку, куда будем сохранять все регистры. Далее из нее сделаем архив, а папку удалим. Если архив существует — удалим его. Считаем, что, если мы зашли в данный метод, то мы уже знаем, что нам нужно генерировать файл (он устарел).

self.registers = [model.__tablename__ for model in RegistersModel.__subclasses__()]
alive_tables_helper = AliveTablesHelper(self.bu, self.version, self.year, self.quarter)
log_info(f"Registers to save: {self.registers}")

Далее определим список регистров и инстанс для AliveTablesHelper.

Затем пробежимся по всем регистрам for register_name in self.registers.

На каждом шаге нужно проверять, не запустился ли расчет. Это более важное событие для нас, и если он запустился, то должен работать именно он, а выгрузка уже устарела.

if alive_tables_helper.check_calc_alive():
	raise 

Далее будем писать текущий статус выгрузки в таблицу.

alive_tables_helper.write_gen_zip_alive_status_to_db(
GeneratedZipSteps.download_register_from_db.value, register_name
)

Если у нас есть months_for_split — значит, у нашей модели в колонках есть месяц и количество строк больше порога. Нужно разбить регистр на отдельные файлы.

month_column_name = model.__table__.columns.get(MONTH_FIELD_NAME).name
for month in months_for_split:
	filters[month_column_name] = str(month)
	xlsx_path = xlsx_files_folder / f"{register_name}_{month}.xlsx"
	df = crud.get_actual_df(
		table_name=register_name,
		bu=bu,
		version=version,
		year=year,
		quarter=quarter,
		filters=filters,
	)
	log_info(f"Register {register_name} was downloaded from DB. Month: {month}.")
	alive_tables_helper.write_gen_zip_alive_status_to_db(
		GeneratedZipSteps.save_register_on_disk.value, register_name
	)
	with ExcelWriter(xlsx_path, engine="xlsxwriter") as writer:
		formats = ExcelFormat(writer.book)
		ReportExcelWriter(
			crud=crud, df=df, db_object_name=register_name, writer=writer, 				formats=formats
		).write()
	log_info(f"Register {register_name} was saved on disk. Month: {month}.")
	del filters[month_column_name]
	del df

Теперь мы делаем все то же самое, что и в изначальной версии, но мы добавили фильтр по месяцу и сохранение статусов в таблицу. Важным моментом будет удалить в конце filters — чтобы он не начал применяться для следующей таблицы, если там записей меньше порога.

xlsx_path = xlsx_files_folder / f"{register_name}.xlsx"

df = crud.get_actual_df(
  table_name=register_name,
  bu=bu,
  version=version,
  year=year,
  quarter=quarter
)
log_info(f"Register {register_name} was downloaded from DB.")
with ExcelWriter(xlsx_path, engine="xlsxwriter") as writer:
	formats = ExcelFormat(writer.book)
	ReportExcelWriter(
		crud=crud, df=df, db_object_name=register_name, writer=writer, formats=formats
	).write()
del df
alive_tables_helper.write_gen_zip_alive_status_to_db(GeneratedZipSteps.save_register_on_disk.value, register_name)
log_info(f"Register {register_name} was saved on disk.")

Для остальных же случаев выгрузка почти никак не будет отличаться от изначального варианта.

Теперь нам остается дописать отдачу файла.

@stream_with_context
def file_streamer(self) -> Iterator[bytes]:
	file_path = self._get_file()
	if not file_path:
		raise FileNotFoundError(Файл не существует.")
	with open(file_path, "rb") as file:
		while True:
			data = file.read(self.chunk_size)
			if not data:
				break
			yield dat

И запуск генерации в отдельном потоке.

def run_thread_for_generate_zip(self) -> None:
	app_context = current_app.app_context()

	build_thread = threading.Thread(
		target=self._generate_zip,
		kwargs={
                "app_context": app_context,
        },
      daemon=True,
    )
    build_thread.start()

Далее остается только собрать все в нашем endpoint, который будет работать по следующему алгоритму:

1. Если запущен расчет — загрузка файла недоступна.

2. Если расчет ни разу не запускался — скачивание файла недоступно.

3. Если архив с текущими параметрами генерируется — скачивание недоступно, но можно пользователю вернуть текущий статус формирования файла.

4. Если файл существует и он актуальный — вернем его пользователю.

5. Проверим, доступно ли нам необходимое место на диске — если места недостаточно, вернем ошибку.

6. Запишем в таблицу БД статус, что формирование файла началось, и начнем создавать новый файл.

Endpoint получается достаточно простой, когда у вас реализованы все необходимые компоненты.

Таким образом, новая схема выгрузки стала следующей:

 Новая схема выгрузки данных
Новая схема выгрузки данных

Cronjob

Решив проблему выгрузки регистров добавлением кэширования архивов на диск, мы задали себе вопрос, на который необходимо было ответить, чтобы выкатить результат пользователю: «Что будет, если дисковое пространство закончится?».

Мы не можем кэшировать архивы и хранить их вечно. Когда будут добавлять новые балансовые единицы, они также будут хранить свои версии архивов на диске, поэтому нужно продумать механизм, который бы выполнял очистку старых архивов. Для этого мы написали cronjob. Она запускается раз в неделю и удаляет архивы, которые лежат на диске дольше всего. Этим решением мы исключаем проблему переполнения диска. Конечно, можно сделать более продвинутое решение и отслеживать те архивы, которые пользуются большей популярностью у пользователей, но на текущий момент тактика, выбранная нами, работает.

Заключение

Этим улучшением нам удалось решить проблему перезагрузки Pod в тестовой системе и сократить время выгрузки данных пользователю. В новом решении выгрузка занимает не более 20 минут вместо 30, мы не держим так долго http‑соединение, а фронт использует так называемый метод ping pong — он запрашивает у бэка статус и анализирует его.

  • Если пришел статус формирования выгрузки — показываем его пользователю.

  • Если пришел файл — отдаем его пользователю.

  • Если вернулась ошибка — покажем ее пользователю, чтобы он решил, что делать дальше.

При этом, если архив с регистрами уже был сформирован, отдадим его пользователю моментально, не затрачивая времени на его создание, как раньше.

Также мы решили вопрос с чрезмерным потреблением RAM. Раньше мы на бэкенд отводили 12 Gb памяти, но этого нам не хватало. Сейчас же в пике мы потребляем до 3.5 Gb. За счет того, что мы решили вопрос потребления памяти, мы также ушли от проблемы перезагрузки Pod.


ДО

После

Первая выгрузка

~30 мин

До 20 мин

Повторная выгрузка

~30 мин

Моментально

Перезагрузка Pod

Почти каждый раз при выгрузке

0

RAM

Не хватало 12 Gb

Менее 3.5 Gb

При этом мы также показываем пользователю статус выгрузки архива и статус‑бар. Пользователь может увидеть, на каком этапе происходит выгрузка, и примерно оценить ее время.

Статусы выгрузки
Статусы выгрузки