/
simon202
/
LAB11_CHAOS_ENGINEERING
Обзор
Документация
Войти
/
simon202
/
LAB11_CHAOS_ENGINEERING
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
cloud_client.py
250 строк
9 KB
simon202
upload files
20 дек 2025, 09:40
20 дек 2025, 09:40
3f2520b
Код
Авторство
О чём код?
""" Cloud Data Client - базовый компонент для работы с AWS сервисами (LocalStack) Взято из Лабораторной работы 10 и адаптировано для Chaos Engineering """ import boto3 import json import pandas as pd from io import StringIO from datetime import datetime class CloudDataClient: def __init__(self, use_localstack=True): """Инициализация клиента для работы с AWS сервисами""" self.use_localstack = use_localstack if use_localstack: # Конфигурация для LocalStack self.endpoint_url = "http://localhost:4566" self.region = "us-east-1" self.aws_access_key_id = "test" self.aws_secret_access_key = "test" else: # Реальные AWS (требуются настоящие credentials) self.endpoint_url = None self.region = "us-east-1" self.aws_access_key_id = None self.aws_secret_access_key = None # Создаем клиенты AWS сервисов self._init_clients() print(f"✅ CloudDataClient инициализирован ({'LocalStack' if use_localstack else 'AWS'})") def _init_clients(self): """Инициализация AWS клиентов""" client_config = { 'region_name': self.region, 'aws_access_key_id': self.aws_access_key_id, 'aws_secret_access_key': self.aws_secret_access_key } if self.endpoint_url: client_config['endpoint_url'] = self.endpoint_url self.s3_client = boto3.client('s3', **client_config) self.sqs_client = boto3.client('sqs', **client_config) # ==================== S3 ОПЕРАЦИИ ==================== def create_bucket(self, bucket_name): """Создает S3 bucket""" try: self.s3_client.create_bucket(Bucket=bucket_name) print(f"📦 Bucket '{bucket_name}' создан") return True except self.s3_client.exceptions.BucketAlreadyOwnedByYou: print(f"📦 Bucket '{bucket_name}' уже существует") return True except Exception as e: print(f"❌ Ошибка создания bucket: {e}") return False def upload_csv_to_s3(self, dataframe, bucket_name, file_key): """Загружает DataFrame как CSV в S3""" try: csv_buffer = StringIO() dataframe.to_csv(csv_buffer, index=False) self.s3_client.put_object( Bucket=bucket_name, Key=file_key, Body=csv_buffer.getvalue() ) print(f"📤 Файл '{file_key}' загружен в bucket '{bucket_name}'") 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""" try: 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}' скачан из bucket '{bucket_name}'") return dataframe except Exception as e: print(f"❌ Ошибка скачивания из S3: {e}") return None def list_bucket_files(self, bucket_name): """Список файлов в bucket""" try: response = self.s3_client.list_objects_v2(Bucket=bucket_name) files = [obj['Key'] for obj in response.get('Contents', [])] return files except Exception as e: print(f"❌ Ошибка получения списка файлов: {e}") return [] def delete_file_from_s3(self, bucket_name, file_key): """Удаляет файл из S3""" try: self.s3_client.delete_object(Bucket=bucket_name, Key=file_key) print(f"🗑️ Файл '{file_key}' удален из bucket '{bucket_name}'") return True except Exception as e: print(f"❌ Ошибка удаления файла: {e}") return False # ==================== SQS ОПЕРАЦИИ ==================== def create_queue(self, queue_name): """Создает SQS очередь""" try: 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): """Отправляет сообщение в очередь""" try: if isinstance(message_body, dict): message_body = json.dumps(message_body) response = self.sqs_client.send_message( QueueUrl=queue_url, MessageBody=message_body ) print(f"📤 Сообщение отправлено в очередь (ID: {response['MessageId']})") return response['MessageId'] except Exception as e: print(f"❌ Ошибка отправки сообщения: {e}") return None def receive_messages(self, queue_url, max_messages=10, wait_time=5): """Получает сообщения из очереди""" try: response = self.sqs_client.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=max_messages, WaitTimeSeconds=wait_time ) messages = [] for msg in response.get('Messages', []): try: body = json.loads(msg['Body']) except json.JSONDecodeError: body = msg['Body'] messages.append({ 'message_id': msg['MessageId'], 'receipt_handle': msg['ReceiptHandle'], 'body': body }) print(f"📥 Получено {len(messages)} сообщений из очереди") return messages except Exception as e: print(f"❌ Ошибка получения сообщений: {e}") return [] def delete_message(self, queue_url, receipt_handle): """Удаляет сообщение из очереди""" try: self.sqs_client.delete_message( QueueUrl=queue_url, ReceiptHandle=receipt_handle ) print("🗑️ Сообщение удалено из очереди") return True except Exception as e: print(f"❌ Ошибка удаления сообщения: {e}") return False def get_queue_attributes(self, queue_url): """Получает атрибуты очереди""" try: response = self.sqs_client.get_queue_attributes( QueueUrl=queue_url, AttributeNames=['All'] ) return response.get('Attributes', {}) except Exception as e: print(f"❌ Ошибка получения атрибутов очереди: {e}") return {} # Пример использования if __name__ == "__main__": # Создаем клиент для LocalStack client = CloudDataClient(use_localstack=True) # Тестируем S3 print("\n🧪 ТЕСТИРОВАНИЕ S3") print("=" * 40) # Создаем bucket client.create_bucket("test-bucket") # Создаем тестовые данные test_data = pd.DataFrame({ 'id': [1, 2, 3], 'name': ['Alice', 'Bob', 'Charlie'], 'salary': [50000, 60000, 70000] }) # Загружаем данные client.upload_csv_to_s3(test_data, "test-bucket", "employees.csv") # Скачиваем данные downloaded_data = client.download_csv_from_s3("test-bucket", "employees.csv") if downloaded_data is not None: print("Скачанные данные:") print(downloaded_data) # Тестируем SQS print("\n🧪 ТЕСТИРОВАНИЕ SQS") print("=" * 40) # Создаем очередь queue_url = client.create_queue("test-queue") if queue_url: # Отправляем сообщение message = {"event": "data_uploaded", "timestamp": datetime.now().isoformat()} client.send_message(queue_url, message) # Получаем сообщения messages = client.receive_messages(queue_url) for msg in messages: print(f"Получено сообщение: {msg['body']}") client.delete_message(queue_url, msg['receipt_handle']) print("\n✅ Тестирование CloudDataClient завершено!")