/
dievavar
/
lab10_cloud_testing
Обзор
Документация
Войти
/
dievavar
/
lab10_cloud_testing
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
cloud_client.py
311 строк
12 KB
dievavar
upload files
06 дек 2025, 11:50
06 дек 2025, 11:50
85ccbcc
Код
Авторство
О чём код?
import boto3 import pandas as pd import json from io import StringIO, BytesIO import os import time class CloudDataClient: def __init__(self, use_localstack=True): self.use_localstack = use_localstack self.setup_clients() def check_localstack_connection(self): """Проверяем подключение к LocalStack""" if not self.use_localstack: return True import requests try: response = requests.get("http://localhost:4566/_localstack/health", timeout=5) if response.status_code == 200: print("✅ LocalStack подключен") return True else: print(f"❌ LocalStack отвечает с кодом {response.status_code}") return False except requests.exceptions.ConnectionError: print("❌ Не могу подключиться к LocalStack. Запустите его сначала!") print(" Используйте: python localstack_manager.py") return False except Exception as e: print(f"❌ Ошибка проверки подключения: {e}") return False def setup_clients(self): """Настраиваем клиенты для AWS сервисов""" if not self.check_localstack_connection(): print("⚠️ Продолжаем без подключения к LocalStack") if self.use_localstack: try: # Используем LocalStack для локального тестирования self.s3_client = boto3.client( 's3', endpoint_url='http://localhost:4566', aws_access_key_id='test', aws_secret_access_key='test', region_name='us-east-1', config=boto3.session.Config( connect_timeout=5, read_timeout=5, retries={'max_attempts': 3} ) ) self.sqs_client = boto3.client( 'sqs', endpoint_url='http://localhost:4566', aws_access_key_id='test', aws_secret_access_key='test', region_name='us-east-1', config=boto3.session.Config( connect_timeout=5, read_timeout=5, retries={'max_attempts': 3} ) ) print("✅ Клиенты AWS настроены для LocalStack") except Exception as e: print(f"⚠️ Ошибка настройки клиентов: {e}") print("⚠️ Создаем заглушки для клиентов") self.s3_client = None self.sqs_client = None else: # Используем реальные AWS сервисы self.s3_client = boto3.client('s3') self.sqs_client = boto3.client('sqs') print("✅ Клиенты AWS настроены для реального окружения") # S3 операции def create_bucket(self, bucket_name): """Создаем S3 bucket""" if not self.s3_client: print("❌ S3 клиент не инициализирован") return False try: print(f"🔄 Создаем bucket '{bucket_name}'...") if self.use_localstack: # Для LocalStack используем простой вызов self.s3_client.create_bucket(Bucket=bucket_name) else: self.s3_client.create_bucket( Bucket=bucket_name, CreateBucketConfiguration={'LocationConstraint': 'us-east-1'} ) print(f"✅ Bucket '{bucket_name}' создан") return True except Exception as e: print(f"❌ Ошибка создания bucket: {e}") print(f" Тип ошибки: {type(e).__name__}") print(f" Проверьте что LocalStack запущен на порту 4566") print(f" Запустите: docker-compose up -d") return False def upload_csv_to_s3(self, dataframe, bucket_name, file_key): """Загружаем DataFrame в S3 как CSV""" if not self.s3_client: print("❌ S3 клиент не инициализирован") return False try: # Конвертируем DataFrame в CSV csv_buffer = StringIO() dataframe.to_csv(csv_buffer, index=False) # Загружаем в S3 print(f"🔄 Загружаем '{file_key}' в bucket '{bucket_name}'...") self.s3_client.put_object( Bucket=bucket_name, Key=file_key, Body=csv_buffer.getvalue() ) print(f"✅ Файл '{file_key}' загружен в S3") return True except Exception as e: print(f"❌ Ошибка загрузки в S3: {e}") return False def download_csv_from_s3(self, bucket_name, file_key): """Скачиваем CSV из S3 и возвращаем DataFrame""" if not self.s3_client: print("❌ S3 клиент не инициализирован") return None try: print(f"🔄 Скачиваем '{file_key}' из bucket '{bucket_name}'...") response = self.s3_client.get_object(Bucket=bucket_name, Key=file_key) csv_content = response['Body'].read().decode('utf-8') dataframe = pd.read_csv(StringIO(csv_content)) print(f"✅ Файл '{file_key}' скачан из S3") return dataframe except Exception as e: print(f"❌ Ошибка скачивания из S3: {e}") return None def list_bucket_files(self, bucket_name): """Получаем список файлов в bucket""" if not self.s3_client: print("❌ S3 клиент не инициализирован") return [] try: print(f"🔄 Получаем список файлов из bucket '{bucket_name}'...") response = self.s3_client.list_objects_v2(Bucket=bucket_name) if 'Contents' in response: files = [obj['Key'] for obj in response['Contents']] print(f"📁 Файлы в bucket '{bucket_name}': {files}") return files else: print(f"📁 Bucket '{bucket_name}' пуст") return [] except Exception as e: print(f"❌ Ошибка получения списка файлов: {e}") return [] # SQS операции def create_queue(self, queue_name): """Создаем SQS очередь""" if not self.sqs_client: print("❌ SQS клиент не инициализирован") return None try: print(f"🔄 Создаем очередь '{queue_name}'...") response = self.sqs_client.create_queue(QueueName=queue_name) queue_url = response['QueueUrl'] print(f"✅ Очередь '{queue_name}' создана: {queue_url}") return queue_url except Exception as e: print(f"❌ Ошибка создания очереди: {e}") return None def send_message(self, queue_url, message_body): """Отправляем сообщение в SQS очередь""" if not self.sqs_client: print("❌ SQS клиент не инициализирован") return None try: print(f"🔄 Отправляем сообщение в очередь...") response = self.sqs_client.send_message( QueueUrl=queue_url, MessageBody=json.dumps(message_body) ) print(f"✅ Сообщение отправлено: {message_body}") return response['MessageId'] except Exception as e: print(f"❌ Ошибка отправки сообщения: {e}") return None def receive_messages(self, queue_url, max_messages=10): """Получаем сообщения из SQS очереди""" if not self.sqs_client: print("❌ SQS клиент не инициализирован") return [] try: print(f"🔄 Получаем сообщения из очереди...") response = self.sqs_client.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=max_messages, WaitTimeSeconds=5 ) messages = [] if 'Messages' in response: for msg in response['Messages']: message_body = json.loads(msg['Body']) messages.append({ 'body': message_body, 'receipt_handle': msg['ReceiptHandle'] }) print(f"✅ Получено {len(messages)} сообщений") else: print("📭 Нет новых сообщений") return messages except Exception as e: print(f"❌ Ошибка получения сообщений: {e}") return [] def delete_message(self, queue_url, receipt_handle): """Удаляем сообщение из очереди""" if not self.sqs_client: print("❌ SQS клиент не инициализирован") return False try: self.sqs_client.delete_message( QueueUrl=queue_url, ReceiptHandle=receipt_handle ) print("✅ Сообщение удалено из очереди") return True except Exception as e: print(f"❌ Ошибка удаления сообщения: {e}") return False # Пример использования с проверкой if __name__ == "__main__": print("=" * 60) print("Тестирование CloudDataClient с LocalStack") print("=" * 60) # Создаем клиент для локального тестирования client = CloudDataClient(use_localstack=True) # Проверяем подключение if not client.check_localstack_connection(): print("\n❌ LocalStack не запущен!") print("Запустите его командой:") print("1. Убедитесь что Docker запущен") print("2. В терминале выполните: docker-compose up -d") print("3. Подождите 10-15 секунд пока LocalStack запустится") print("4. Повторите запуск теста") exit(1) print("\n" + "=" * 60) print("Тестируем S3 операции") print("=" * 60) # Тестируем S3 bucket_created = client.create_bucket("test-bucket") if bucket_created: # Создаем тестовые данные test_data = pd.DataFrame({ 'id': [1, 2, 3], 'name': ['Alice', 'Bob', 'Charlie'], 'value': [100, 200, 300] }) client.upload_csv_to_s3(test_data, "test-bucket", "test-data.csv") files = client.list_bucket_files("test-bucket") if files: downloaded_data = client.download_csv_from_s3("test-bucket", "test-data.csv") if downloaded_data is not None: print("\n📊 Скачанные данные:") print(downloaded_data) print("\n" + "=" * 60) print("Тестируем SQS операции") print("=" * 60) # Тестируем SQS queue_url = client.create_queue("test-queue") if queue_url: client.send_message(queue_url, {"type": "test", "data": "hello"}) messages = client.receive_messages(queue_url) for msg in messages: print(f"📨 Получено сообщение: {msg['body']}") client.delete_message(queue_url, msg['receipt_handle']) print("\n" + "=" * 60) print("Тестирование завершено!") print("=" * 60)