/
Demek
/
DSam621
Обзор
Документация
Войти
/
Demek
/
DSam621
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
Python/excel-to-bd/db_manager.py
374 строки
15 KB
d_e_m_e_k
обнвление проекто, но ошибка все еще есть
07 авг 2026, 09:18
07 авг 2026, 09:18
0a6239a
Код
Авторство
О чём код?
import psycopg2 from psycopg2 import sql, extras from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT import pandas as pd from sqlalchemy import create_engine from config import Config import logging logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class DatabaseManager: def __init__(self): self.db_url = Config.get_db_url() self.connection = None self.engine = None self.db_name = Config.DB_NAME self.db_user = Config.DB_USER self.db_password = Config.DB_PASSWORD self.db_host = Config.DB_HOST self.db_port = Config.DB_PORT def check_database_exists(self): """Проверяет существование базы данных""" try: # Подключаемся к стандартной базе postgres для проверки conn = psycopg2.connect( host=self.db_host, port=self.db_port, database='postgres', user=self.db_user, password=self.db_password ) conn.autocommit = True with conn.cursor() as cursor: # Проверяем существование БД cursor.execute( "SELECT 1 FROM pg_database WHERE datname = %s", (self.db_name,) ) exists = cursor.fetchone() is not None conn.close() if exists: logger.info(f"✅ База данных '{self.db_name}' существует") else: logger.info(f"ℹ️ База данных '{self.db_name}' не найдена") return exists except psycopg2.OperationalError as e: if "does not exist" in str(e): logger.info(f"ℹ️ База данных '{self.db_name}' не существует") return False else: logger.error(f"❌ Ошибка подключения: {e}") raise def create_database(self): """Создает базу данных если она не существует""" try: # Подключаемся к стандартной базе postgres conn = psycopg2.connect( host=self.db_host, port=self.db_port, database='postgres', user=self.db_user, password=self.db_password ) conn.autocommit = True with conn.cursor() as cursor: # Создаем базу данных cursor.execute( sql.SQL("CREATE DATABASE {}").format( sql.Identifier(self.db_name) ) ) logger.info(f"✅ База данных '{self.db_name}' успешно создана") conn.close() return True except psycopg2.OperationalError as e: logger.error(f"❌ Ошибка создания базы данных: {e}") return False except Exception as e: logger.error(f"❌ Непредвиденная ошибка при создании БД: {e}") return False def ensure_database(self): """Гарантирует существование базы данных""" if not self.check_database_exists(): logger.info(f"🔄 Создаем базу данных '{self.db_name}'...") if self.create_database(): logger.info(f"✅ База данных '{self.db_name}' создана") return True else: logger.error(f"❌ Не удалось создать базу данных '{self.db_name}'") return False return True def connect(self): """Устанавливает соединение с БД""" try: # Сначала проверяем/создаем БД if not self.ensure_database(): return False # Теперь подключаемся к целевой БД self.connection = psycopg2.connect( host=self.db_host, port=self.db_port, database=self.db_name, user=self.db_user, password=self.db_password ) # Создаем SQLAlchemy engine self.engine = create_engine(self.db_url) logger.info(f"✅ Подключение к PostgreSQL установлено (БД: {self.db_name})") return True except psycopg2.OperationalError as e: logger.error(f"❌ Ошибка подключения: {e}") return False except Exception as e: logger.error(f"❌ Непредвиденная ошибка: {e}") return False def create_table(self, table_name=None, schema=None): """Создает таблицу для данных""" if table_name is None: table_name = Config.TABLE_NAME # Если схема не передана, используем стандартную if schema is None: schema = """ id SERIAL PRIMARY KEY, record_id INTEGER, date DATE, customer_name VARCHAR(200), customer_email VARCHAR(200), product_category VARCHAR(100), product_name VARCHAR(200), quantity INTEGER, price DECIMAL(10, 2), discount DECIMAL(5, 2), is_active BOOLEAN, rating DECIMAL(3, 1), notes TEXT, total DECIMAL(12, 2), created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP """ create_table_query = f""" CREATE TABLE IF NOT EXISTS {table_name} ( {schema} ); -- Индексы для ускорения запросов CREATE INDEX IF NOT EXISTS idx_{table_name}_date ON {table_name}(date); CREATE INDEX IF NOT EXISTS idx_{table_name}_category ON {table_name}(product_category); CREATE INDEX IF NOT EXISTS idx_{table_name}_customer ON {table_name}(customer_name); -- Триггер для обновления updated_at CREATE OR REPLACE FUNCTION update_updated_at_column() RETURNS TRIGGER AS $$ BEGIN NEW.updated_at = CURRENT_TIMESTAMP; RETURN NEW; END; $$ language 'plpgsql'; DROP TRIGGER IF EXISTS update_{table_name}_updated_at ON {table_name}; CREATE TRIGGER update_{table_name}_updated_at BEFORE UPDATE ON {table_name} FOR EACH ROW EXECUTE FUNCTION update_updated_at_column(); """ try: with self.connection.cursor() as cursor: cursor.execute(create_table_query) self.connection.commit() logger.info(f"✅ Таблица {table_name} создана/обновлена") return True except Exception as e: logger.error(f"❌ Ошибка создания таблицы: {e}") self.connection.rollback() return False def table_exists(self, table_name): """Проверяет существование таблицы""" query = """ SELECT EXISTS ( SELECT 1 FROM information_schema.tables WHERE table_name = %s ); """ try: with self.connection.cursor() as cursor: cursor.execute(query, (table_name,)) exists = cursor.fetchone()[0] return exists except Exception as e: logger.error(f"❌ Ошибка проверки таблицы: {e}") return False def load_data_from_dataframe(self, df, table_name=None): """Загружает данные из DataFrame в PostgreSQL""" if table_name is None: table_name = Config.TABLE_NAME try: # Очищаем данные перед загрузкой df_clean = df.copy() # Преобразуем типы данных for col in df_clean.columns: if 'date' in col.lower() or 'дата' in col.lower(): try: df_clean[col] = pd.to_datetime(df_clean[col]).dt.date except: pass if 'is_active' in col.lower() or 'статус' in col.lower(): # Преобразуем булевы значения if df_clean[col].dtype == 'object': df_clean[col] = df_clean[col].map({ 'Активный': True, 'Неактивный': False, 'True': True, 'False': False, 'Yes': True, 'No': False }).fillna(False) # Обрабатываем пропуски df_clean = df_clean.where(pd.notnull(df_clean), None) # Проверяем существование таблицы if not self.table_exists(table_name): logger.warning(f"⚠️ Таблица {table_name} не существует, создаем...") if not self.create_table(table_name): logger.error("❌ Не удалось создать таблицу") return False # Используем SQLAlchemy для массовой загрузки df_clean.to_sql( table_name, self.engine, if_exists='append', index=False, method='multi', chunksize=Config.CHUNK_SIZE ) logger.info(f"✅ Загружено {len(df_clean)} записей в таблицу {table_name}") return True except Exception as e: logger.error(f"❌ Ошибка загрузки данных: {e}") return False def load_data_from_excel(self, excel_file, table_name=None): """Загружает данные из Excel в БД""" if table_name is None: table_name = Config.TABLE_NAME try: # Читаем Excel df = pd.read_excel(excel_file) return self.load_data_from_dataframe(df, table_name) except Exception as e: logger.error(f"❌ Ошибка чтения Excel: {e}") return False def execute_query(self, query, params=None): """Выполняет произвольный запрос""" try: with self.connection.cursor() as cursor: if params: cursor.execute(query, params) else: cursor.execute(query) if cursor.description: # Если это SELECT columns = [desc[0] for desc in cursor.description] data = cursor.fetchall() return pd.DataFrame(data, columns=columns) else: # Если это INSERT/UPDATE/DELETE self.connection.commit() return f"✅ Выполнено. Затронуто строк: {cursor.rowcount}" except Exception as e: logger.error(f"❌ Ошибка выполнения запроса: {e}") self.connection.rollback() return None def get_table_info(self, table_name=None): """Получает информацию о таблице""" if table_name is None: table_name = Config.TABLE_NAME query = """ SELECT column_name, data_type, is_nullable, column_default FROM information_schema.columns WHERE table_name = %s ORDER BY ordinal_position; """ return self.execute_query(query, (table_name,)) def get_statistics(self, table_name=None): """Получает статистику по данным""" if table_name is None: table_name = Config.TABLE_NAME query = f""" SELECT COUNT(*) as total_records, COUNT(DISTINCT customer_name) as unique_customers, AVG(price) as avg_price, SUM(total) as total_sum, MIN(date) as first_date, MAX(date) as last_date, COUNT(DISTINCT product_category) as categories_count FROM {table_name}; """ return self.execute_query(query) def get_table_list(self): """Получает список всех таблиц в БД""" query = """ SELECT table_name FROM information_schema.tables WHERE table_schema = 'public' ORDER BY table_name; """ return self.execute_query(query) def drop_table(self, table_name=None, confirm=False): """Удаляет таблицу (с подтверждением)""" if table_name is None: table_name = Config.TABLE_NAME if not confirm: logger.warning(f"⚠️ Для удаления таблицы {table_name} установите confirm=True") return False try: with self.connection.cursor() as cursor: cursor.execute(f"DROP TABLE IF EXISTS {table_name} CASCADE;") self.connection.commit() logger.info(f"✅ Таблица {table_name} удалена") return True except Exception as e: logger.error(f"❌ Ошибка удаления таблицы: {e}") self.connection.rollback() return False def close(self): """Закрывает соединение""" if self.connection: self.connection.close() logger.info("🔒 Соединение закрыто")