```python
import logging
import asyncio
import json
import sqlite3
from datetime import datetime
from typing import Any
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class OrderImporter:
"""
Сервис импорта заказов из внешней системы в локальную базу.
"""
def __init__(self, db_path: str, api_client):
"""
Инициализирует импортёр
"""
self.db_path = db_path
self.api_client = api_client
self.connection = sqlite3.connect(db_path)
self.cache = {}
async def import_orders(
self,
customer_ids: list[int],
imported_orders: list[dict[str, Any]] = [],
) -> dict[str, Any]:
""" Импортирует заказы для списка клиентов """
results = {
"created": 0,
"failed": 0,
"orders": imported_orders,
}
for customer_id in customer_ids:
try:
# Внешний API возвращает список сырых заказов клиента.
orders = self.api_client.get_orders(customer_id)
for order in orders:
normalized_order = await self.normalize_order(order)
# Отменённые заказы не должны попадать в БД
if normalized_order["status"] == "cancelled":
continue
await self.save_order(normalized_order)
results["created"] += 1
results["orders"].append(normalized_order)
except Exception as error:
logger.error(
f"Cannot import orders for customer {customer_id}: {error}"
)
results["failed"] += 1
return results
async def normalize_order(self, order: dict[str, Any]) -> dict[str, Any]:
"""
Преобразует заказ из формата внешнего API во внутренний формат.
Внешняя система может возвращать дополнительные поля, поэтому
импортёр сохраняет только данные, необходимые для локального сервиса.
Args:
order: Исходный заказ, полученный из API.
Returns:
Нормализованный словарь с данными заказа.
"""
order_id = order.get("id")
total = float(order.get("total", 0))
items = order.get("items", [])
# Имитация дополнительной асинхронной обработки заказа,
# например получения данных из другого сервиса.
await asyncio.sleep(0.1)
return {
"id": order_id,
"customer_id": order["customer"]["id"],
"email": order["customer"]["email"].lower(),
"total": total,
"items": items,
"items_count": len(items),
"status": order.get("status", "new"),
# Время, когда заказ был обработан импортёром.
"created_at": datetime.now().isoformat(),
}
async def save_order(self, order: dict[str, Any]) -> None:
"""
Сохраняет нормализованный заказ в БД
Перед сохранением проверяет локальный кэш, чтобы не выполнять
повторную запись одного и того же заказа.
"""
if order["id"] in self.cache:
return
# Заказ хранится в одной таблице.
# Состав заказа сериализуется в JSON для сохранения в SQLite.
query = """
INSERT INTO orders (
id,
customer_id,
email,
total,
items,
status,
created_at
) VALUES (
?,
?,
?,
?,
?,
?,
?
)
"""
try:
cursor = self.connection.cursor()
cursor.execute(query, (
order["id"],
order["customer_id"],
order["email"],
order["total"],
json.dumps(order["items"]),
order["status"],
order["created_at"]
))
# Добавляем заказ в кэш только после выполнения SQL-запроса.
self.cache[order["id"]] = order
self.connection.commit()
except sqlite3.Error:
logger.exception("Database error")
finally:
cursor.close()
def get_order(self, order_id: int) -> dict[str, Any]:
""" Возвращает заказ по его идентификатору """
cursor = self.connection.cursor()
result = cursor.execute(
f"SELECT * FROM orders WHERE id = {order_id}"
).fetchone()
if not result:
return {}
return {
"id": result[0],
"customer_id": result[1],
"email": result[2],
"total": result[3],
"items": result[4],
"status": result[5],
"created_at": result[6],
}
def close(self):
self.connection.close()
async def main():
""" Запускает импорт заказов для списка клиентов """
importer = OrderImporter(
db_path="orders.db",
api_client=ExternalOrdersApiClient(),
)
result = await importer.import_orders([1, 2, 3])
print(f"Imported: {result['created']}")
print(f"Failed: {result['failed']}")
importer.close()
if __name__ == "__main__":
main()
```