/
nanocubit
/
ZeroLink
Обзор
Документация
Войти
/
nanocubit
/
ZeroLink
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
examples/optimized_protocol.py
178 строк
7 KB
slavam
feat: добавлена рабочая реализация GPU пула памяти, обновлены все компоненты ZeroLink
23 янв 2026, 20:56
23 янв 2026, 20:56
0106973
Код
Авторство
О чём код?
""" Пример оптимизированного протокола для ZeroLink Этот файл демонстрирует потенциальные оптимизации для zerolink/core/protocol/messages.py """ import struct import io from typing import Any, Dict, List import pickle import numpy as np from memory_profiler import profile class OptimizedProtocol: """ Оптимизированный протокол с использованием memoryview и буферов """ # Форматы сообщений MSG_HDR_FMT = "<IIQ" # magic, version, length MSG_HDR_SIZE = struct.calcsize(MSG_HDR_FMT) def __init__(self): # Буферы для уменьшения выделения памяти self.header_buffer = bytearray(self.MSG_HDR_SIZE) self.payload_buffer = bytearray(1024 * 1024) # 1MB буфер def pack_message_optimized(self, msg_type: int, payload: bytes) -> bytes: """ Оптимизированная упаковка сообщения с использованием memoryview """ payload_len = len(payload) # Заполняем заголовок header_view = memoryview(self.header_buffer)[:self.MSG_HDR_SIZE] struct.pack_into(self.MSG_HDR_FMT, header_view, 0, 0xDEADBEEF, 1, payload_len) # Создаем общий буфер для заголовка и полезной нагрузки message_buffer = bytearray(self.MSG_HDR_SIZE + payload_len) message_view = memoryview(message_buffer) # Копируем заголовок message_view[:self.MSG_HDR_SIZE] = header_view # Копируем полезную нагрузку message_view[self.MSG_HDR_SIZE:] = payload return bytes(message_buffer) def unpack_message_optimized(self, data: bytes) -> tuple: """ Оптимизированная распаковка сообщения с использованием memoryview """ data_mv = memoryview(data) # Распаковываем заголовок magic, version, payload_len = struct.unpack_from(self.MSG_HDR_FMT, data_mv) if magic != 0xDEADBEEF: raise ValueError("Неверный магический номер") if len(data_mv) < self.MSG_HDR_SIZE + payload_len: raise ValueError("Недостаточно данных для распаковки") # Возвращаем view на полезную нагрузку без копирования payload_view = data_mv[self.MSG_HDR_SIZE:self.MSG_HDR_SIZE + payload_len] return version, payload_view.tobytes() def serialize_tensor_optimized(self, tensor) -> bytes: """ Оптимизированная сериализация тензора """ # Используем numpy для более быстрой сериализации if hasattr(tensor, 'numpy'): np_array = tensor.numpy() else: np_array = np.asarray(tensor) # Используем BytesIO для создания буфера в памяти buffer = io.BytesIO() np.save(buffer, np_array, allow_pickle=False) return buffer.getvalue() def deserialize_tensor_optimized(self, data: bytes): """ Оптимизированная десериализация тензора """ buffer = io.BytesIO(data) np_array = np.load(buffer, allow_pickle=False) # Преобразуем обратно в тензор import torch return torch.from_numpy(np_array) class BatchMessageProcessor: """ Обработчик пакетных сообщений для уменьшения накладных расходов """ def __init__(self, max_batch_size: int = 10): self.max_batch_size = max_batch_size self.message_queue = [] self.protocol = OptimizedProtocol() def add_message(self, msg_type: int, payload: bytes): """Добавление сообщения в очередь""" self.message_queue.append((msg_type, payload)) if len(self.message_queue) >= self.max_batch_size: return self.process_batch() return None def process_batch(self) -> bytes: """Обработка пакета сообщений""" if not self.message_queue: return b"" # Сериализуем все сообщения в один пакет batch_data = [] for msg_type, payload in self.message_queue: packed_msg = self.protocol.pack_message_optimized(msg_type, payload) batch_data.append(packed_msg) # Объединяем все сообщения в один буфер total_size = sum(len(msg) for msg in batch_data) combined_buffer = bytearray(total_size) offset = 0 for msg in batch_data: combined_buffer[offset:offset+len(msg)] = msg offset += len(msg) # Очищаем очередь self.message_queue.clear() return bytes(combined_buffer) # Пример использования def example_usage(): """Пример использования оптимизированного протокола""" protocol = OptimizedProtocol() # Создаем тестовое сообщение test_payload = b"Hello, ZeroLink!" * 1000 # 13KB # Упаковка сообщения packed = protocol.pack_message_optimized(1, test_payload) print(f"Упакованное сообщение: {len(packed)} bytes") # Распаковка сообщения version, unpacked_payload = protocol.unpack_message_optimized(packed) print(f"Распакованное сообщение: {len(unpacked_payload)} bytes") # Проверка целостности assert test_payload == unpacked_payload print("✓ Проверка целостности пройдена") # Пример пакетной обработки batch_processor = BatchMessageProcessor(max_batch_size=5) for i in range(7): # 7 сообщений result = batch_processor.add_message(i, f"Message {i}".encode()) if result: print(f"Обработан пакет из {len(batch_processor.message_queue)+1} сообщений") # Обработка оставшихся сообщений remaining_batch = batch_processor.process_batch() if remaining_batch: print(f"Обработан остаточный пакет: {len(remaining_batch)} bytes") if __name__ == "__main__": example_usage()