/
kamball
/
lab11
Обзор
Документация
Войти
/
kamball
/
lab11
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
cloud_client.py
152 строки
6 KB
kambal23123
first_commit
20 дек 2025, 14:14
20 дек 2025, 14:14
15d96ee
Код
Авторство
О чём код?
import boto3 import pandas as pd import io from botocore.exceptions import ClientError class CloudDataClient: def __init__(self, use_localstack=True): self.use_localstack = use_localstack if use_localstack: self.endpoint_url = "http://localhost:4566" self.aws_access_key_id = "test" self.aws_secret_access_key = "test" self.region_name = "us-east-1" else: self.endpoint_url = None self.aws_access_key_id = None self.aws_secret_access_key = None self.region_name = "us-east-1" # Инициализируем клиенты self.s3_client = boto3.client( 's3', endpoint_url=self.endpoint_url, aws_access_key_id=self.aws_access_key_id, aws_secret_access_key=self.aws_secret_access_key, region_name=self.region_name ) self.sqs_client = boto3.client( 'sqs', endpoint_url=self.endpoint_url, aws_access_key_id=self.aws_access_key_id, aws_secret_access_key=self.aws_secret_access_key, region_name=self.region_name ) def create_bucket(self, bucket_name): """Создает S3 bucket""" try: if self.use_localstack: self.s3_client.create_bucket(Bucket=bucket_name) else: self.s3_client.create_bucket( Bucket=bucket_name, CreateBucketConfiguration={'LocationConstraint': self.region_name} ) print(f"✅ Bucket {bucket_name} создан") return True except ClientError as e: if e.response['Error']['Code'] == 'BucketAlreadyOwnedByYou': print(f"ℹ️ Bucket {bucket_name} уже существует") return True print(f"❌ Ошибка создания bucket: {e}") return False def upload_csv_to_s3(self, dataframe, bucket_name, file_key): """Загружает DataFrame в S3 как CSV""" try: csv_buffer = io.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_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(io.StringIO(csv_content)) print(f"✅ Файл {file_key} скачан из {bucket_name}") return dataframe except Exception as e: print(f"❌ Ошибка скачивания из S3: {e}") return None 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): """Отправляет сообщение в SQS очередь""" try: if isinstance(message_body, dict): import json message_body = json.dumps(message_body) response = self.sqs_client.send_message( QueueUrl=queue_url, MessageBody=message_body ) print(f"✅ Сообщение отправлено в очередь") return response except Exception as e: print(f"❌ Ошибка отправки сообщения: {e}") return None def receive_messages(self, queue_url, max_messages=10): """Получает сообщения из SQS очереди""" try: response = self.sqs_client.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=max_messages ) messages = [] if 'Messages' in response: for msg in response['Messages']: try: import json body = json.loads(msg['Body']) except: body = msg['Body'] messages.append({ 'body': body, 'receipt_handle': msg['ReceiptHandle'] }) 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 ) return True except Exception as e: print(f"❌ Ошибка удаления сообщения: {e}") return False