/
Undeadguy
/
Tape_python
Обзор
Документация
Войти
/
Undeadguy
/
Tape_python
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
serializing_manager.py
274 строки
10 KB
Undeadguy
1.0
15 апр 2026, 22:40
15 апр 2026, 22:40
556ac2a
Код
Авторство
О чём код?
import asyncio import base64 import json from asyncio import Lock from asyncio import Queue as AsyncQueue from pathlib import Path from typing import IO, List, Optional, Tuple import file_funcs PYTHON_OVERHEAD_COEFF = 1.4 HEX_SERIALIZATION_KOEFF = 1.3 class SerManager: ram_free: Optional[int] = None file_write_lock: Lock = Lock() bytes_in_usage_lock: Lock = Lock() output: IO files: List[Tuple[Path, int]] # (путь, размер_влияние) ser_type: str max_concurrent: int = 5 # максимальное количество параллельных сериализаций def __init__( self, write_to: IO, files: List[Tuple[Path, int]], ram_limit: Optional[int], serialization: str, ): if ram_limit: self.ram_free = ram_limit else: self.ram_free = float("inf") # без ограничений self.output = write_to # Сортируем по размеру (большие сначала - лучше заполняют память) self.files = sorted(files, key=lambda x: x[1], reverse=True) self.ser_type = serialization self.active_tasks = 0 self.completed = 0 self.total_files = len(files) self.task_queue = AsyncQueue() self.ram_available = self.ram_free self.ram_lock = Lock() async def serialize_all(self): """Главный дирижёр - распределяет задачи и следит за памятью""" # Заполняем очередь задачами for file_path, weight in self.files: await self.task_queue.put((file_path, weight)) # Запускаем воркеров workers = [] for i in range(self.max_concurrent): worker = asyncio.create_task(self._worker(i)) workers.append(worker) # Ждём завершения всех задач await self.task_queue.join() # Останавливаем воркеров for _ in workers: await self.task_queue.put(None) # сигнал остановки await asyncio.gather(*workers) # print(f"[Дирижёр] Все {self.completed} файлов обработаны!") async def _worker(self, worker_id: int): while True: task = await self.task_queue.get() if task is None: self.task_queue.task_done() break file_path, weight = task await self._wait_for_memory(weight) async with self.ram_lock: self.ram_available -= weight # print( # f"[Воркер {worker_id}] Зарезервировал {weight} байт. Свободно: {self.ram_available}" # ) try: # Определяем, первый ли это файл is_first = self.completed == 0 await self.serialize_file( self_id=worker_id, file_path=file_path, ram_limit=self.ram_free if isinstance(self.ram_free, int) else None, output_file=self.output, ser_type=self.ser_type, write_file_lock=self.file_write_lock, is_first_file=is_first, # ← передаём флаг ) self.completed += 1 # print( # f"[Воркер {worker_id}] Завершил {file_path.name} ({self.completed}/{self.total_files})" # ) # except Exception as e: # print(f"[Воркер {worker_id}] Ошибка при обработке {file_path.name}: {e}") finally: async with self.ram_lock: self.ram_available += weight # print( # f"[Воркер {worker_id}] Освободил {weight} байт. Свободно: {self.ram_available}" # ) self.task_queue.task_done() async def _wait_for_memory(self, needed_memory: int): if not isinstance(self.ram_free, int): return # без ограничений while True: async with self.ram_lock: if self.ram_available >= needed_memory: return free = self.ram_available # Ждём немного перед следующей проверкой await asyncio.sleep(0.1) # print(f"[Дирижёр] Ожидание памяти: нужно {needed_memory}, свободно {free}") @staticmethod async def serialize_file( self_id: int, file_path: Path, ram_limit: Optional[int], output_file: IO, ser_type: str, write_file_lock: Lock, is_first_file: bool = True, # добавляем параметр ): file_type_txt, storaging_type = file_funcs.FILE_TYPE_MAP[file_path.suffix] # === ПОЛНЫЙ ФАЙЛ (без ограничения памяти) === if ram_limit is None: if storaging_type == file_funcs.FileStorageType.Text: body = file_path.read_text() else: # Binary body = base64.b64encode(file_path.read_bytes()).decode() new_file = file_funcs.File( file_path.name, str(file_path), file_type_txt, (storaging_type is file_funcs.FileStorageType.Base64), body, ) async with write_file_lock: if ser_type.upper() == "JSONL": output_file.write(json.dumps(new_file.__dict__, ensure_ascii=False) + "\n") elif ser_type.upper() == "TOML": # Разделитель между файлами (кроме первого) if not is_first_file: output_file.write("\n") # [[files]] заголовок output_file.write("[[files]]\n") output_file.write(f'name = "{new_file.name}"\n') output_file.write(f'path = "{new_file.path}"\n') output_file.write(f'file_type = "{new_file.file_type}"\n') output_file.write(f"is_binary = {str(new_file.is_binary).lower()}\n") # Многострочная строка для body output_file.write('body = """\n') output_file.write(new_file.body) output_file.write('\n"""\n') return # === ЧАНКОВЫЙ РЕЖИМ === # Расчёт размера чанка BYTES_PER_CHAR = 2 BASE64_OVERHEAD = 1.33 if storaging_type == file_funcs.FileStorageType.Text: chunk_limit = max(1024, int(ram_limit // BYTES_PER_CHAR)) else: chunk_limit = max(1024, int(ram_limit // BASE64_OVERHEAD)) async with write_file_lock: if ser_type.upper() == "JSONL": output_file.write( json.dumps( { "name": file_path.name, "path": str(file_path), "type": file_type_txt, "is_binary": (storaging_type is file_funcs.FileStorageType.Base64), "chunked": True, } ) + "\n" ) elif ser_type.upper() == "TOML": write_toml_header(output_file, file_path, file_type_txt, storaging_type) chunk_index = 0 if storaging_type == file_funcs.FileStorageType.Text: with open(file_path, "r", encoding="utf-8") as source: while True: chunk_text = source.read(chunk_limit) if not chunk_text: break async with write_file_lock: if ser_type.upper() == "JSONL": output_file.write( json.dumps({"chunk": chunk_index, "body": chunk_text}) + "\n" ) elif ser_type.upper() == "TOML": if chunk_index == 0: output_file.write('body = """\n') output_file.write(chunk_text) chunk_index += 1 else: with open(file_path, "rb") as source: while True: chunk_bytes = source.read(chunk_limit) if not chunk_bytes: break chunk_b64 = base64.b64encode(chunk_bytes).decode("ascii") async with write_file_lock: if ser_type.upper() == "JSONL": output_file.write( json.dumps( { "chunk": chunk_index, "body": chunk_b64, "size": len(chunk_bytes), } ) + "\n" ) elif ser_type.upper() == "TOML": if chunk_index == 0: output_file.write('body = """\n') output_file.write(chunk_b64) chunk_index += 1 if ser_type.upper() == "TOML": async with write_file_lock: output_file.write('\n"""\n') output_file.write(f"total_chunks = {chunk_index}\n") def write_toml_header(output: IO, file_path: Path, file_type: str, storaging_type): """Пишет заголовок для TOML (открывает многострочную строку)""" output.write("[[files]]\n") output.write(f'name = "{file_path.name}"\n') output.write(f'path = "{file_path}"\n') output.write(f'file_type = "{file_type}"\n') output.write(f"is_binary = {storaging_type is file_funcs.FileStorageType.Base64}\n") output.write("chunked = true\n") output.write('body = """\n')