Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 28 additions & 3 deletions docker-compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -153,13 +153,17 @@ services:
build:
context: .
dockerfile: docker/services/recommender.dockerfile
command: python app/entry/main.py
command: python main.py server
depends_on:
- recommender_db
- kafka
environment:
DATABASE_URL: postgresql://rasadov:123456@recommender_db/rasadov?sslmode=disable
KAFKA_BOOTSTRAP_SERVERS: kafka:9092
ARTIFACTS_LOCAL_PATH: /app/artifacts
GRPC_PORT: 8080
volumes:
- recommender_artifacts:/app/artifacts
restart: on-failure

recommender-sync:
Expand All @@ -168,13 +172,32 @@ services:
build:
context: .
dockerfile: docker/services/recommender.dockerfile
command: python app/entry/sync.py
command: python main.py consumer
depends_on:
- recommender_db
- kafka
environment:
DATABASE_URL: postgresql://rasadov:123456@recommender_db/rasadov?sslmode=disable
KAFKA_BOOTSTRAP_SERVERS: kafka:9092
PRODUCT_API: http://product:8080/products
PRODUCT_EVENTS_TOPIC: product_events
INTERACTION_EVENTS_TOPIC: interaction_events
restart: on-failure

recommender-train:
container_name: "recommender-train"
image: "rasadov/recommender-train:latest"
build:
context: .
dockerfile: docker/services/recommender.dockerfile
command: sh -c "while true; do python main.py train; sleep 43200; done"
depends_on:
- recommender_db
environment:
DATABASE_URL: postgresql://rasadov:123456@recommender_db/rasadov?sslmode=disable
ARTIFACTS_LOCAL_PATH: /app/artifacts
volumes:
- recommender_artifacts:/app/artifacts
restart: on-failure

graphql:
Expand Down Expand Up @@ -206,4 +229,6 @@ volumes:
order_db_data:
payment_db_data:
recommender_db_data:
kafka-volume:
recommender_artifacts:
kafka-volume:
zookeeper-volume:
24 changes: 20 additions & 4 deletions docker/services/recommender.dockerfile
Original file line number Diff line number Diff line change
@@ -1,12 +1,28 @@
FROM python:3.11-slim

COPY --from=ghcr.io/astral-sh/uv:latest /uv /uvx /bin/

WORKDIR /app

RUN apt-get update && apt-get install -y build-essential libpq-dev && rm -rf /var/lib/apt/lists/*
RUN apt-get update && apt-get install -y --no-install-recommends \
build-essential \
libpq-dev \
&& rm -rf /var/lib/apt/lists/*

# Install dependencies first so this layer is cached when only app code changes.
COPY recommender/pyproject.toml recommender/uv.lock ./
ENV UV_COMPILE_BYTECODE=1 \
UV_LINK_MODE=copy \
UV_PYTHON_DOWNLOADS=never

COPY recommender/requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
RUN uv sync --frozen --no-install-project --no-dev

COPY recommender /app
RUN uv sync --frozen --no-dev

ENV PATH="/app/.venv/bin:${PATH}" \
PYTHONPATH="/app/app:${PYTHONPATH}"

ENV PYTHONPATH="/app:${PYTHONPATH}"
EXPOSE 8080

WORKDIR /app
Empty file removed recommender/app/db/__init__.py
Empty file.
Empty file removed recommender/app/entry/__init__.py
Empty file.
63 changes: 63 additions & 0 deletions recommender/app/entry/grpc_server.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import grpc
from concurrent import futures
from loguru import logger

from generated.pb import recommender_pb2, recommender_pb2_grpc
from recommendations.load import load_service
from shared.config.settings import GRPC_PORT
from shared.db.repo import get_products_by_ids
from shared.db.session import get_session


class RecommenderServiceServicer(recommender_pb2_grpc.RecommenderServiceServicer):
def __init__(self, service):
self._service = service

def GetRecommendations(self, request, context):
try:
product_ids = self._service.for_user(
user_id=request.user_id,
skip=request.skip or 0,
take=request.take or 5,
)
return _to_response(product_ids)
except Exception as exc:
logger.exception("Failed to get recommendations for user {}", request.user_id)
context.set_code(grpc.StatusCode.INTERNAL)
context.set_details(str(exc))
return recommender_pb2.RecommendationResponse()

def GetRecommendationsBasedOnViewed(self, request, context):
try:
product_ids = self._service.for_viewed_products(
product_ids=list(request.ids),
skip=request.skip or 0,
take=request.take or 5,
)
return _to_response(product_ids)
except Exception as exc:
logger.exception("Failed to get viewed-based recommendations")
context.set_code(grpc.StatusCode.INTERNAL)
context.set_details(str(exc))
return recommender_pb2.RecommendationResponse()


def _to_response(product_ids: list[str]) -> recommender_pb2.RecommendationResponse:
with get_session() as session:
products = get_products_by_ids(session, product_ids)
return recommender_pb2.RecommendationResponse(
recommended_products=[p.to_grpc_model() for p in products]
)


def serve() -> None:
service = load_service()
server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
recommender_pb2_grpc.add_RecommenderServiceServicer_to_server(
RecommenderServiceServicer(service),
server,
)
server.add_insecure_port(f"[::]:{GRPC_PORT}")
logger.info("gRPC server started on port {}", GRPC_PORT)
server.start()
server.wait_for_termination()
74 changes: 74 additions & 0 deletions recommender/app/entry/kafka_consumer.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import json

from kafka import KafkaConsumer
import requests
from loguru import logger

from shared.config.settings import (
INTERACTION_EVENTS_TOPIC,
KAFKA_SERVER,
PRODUCT_EVENTS_TOPIC,
)
from shared.db.repo import (
create_or_update_product,
create_product,
delete_product_by_id,
get_product_by_id,
record_interaction,
)
from shared.db.session import get_session
from shared.kafka.utils import product_is_created_or_updated, product_is_deleted
from shared.product.utils import fetch_product_by_id


def start_kafka_consumer() -> None:
consumer = KafkaConsumer(
PRODUCT_EVENTS_TOPIC,
INTERACTION_EVENTS_TOPIC,
bootstrap_servers=KAFKA_SERVER,
group_id="recommender-sync",
)

for message in consumer:
event = json.loads(message.value)

if message.topic == PRODUCT_EVENTS_TOPIC:
_handle_product_event(event)
elif message.topic == INTERACTION_EVENTS_TOPIC:
_handle_interaction_event(event)


def _handle_product_event(event: dict) -> None:
with get_session() as session:
if product_is_created_or_updated(event):
product_data = event["data"]
logger.info(
"Processing product event {} for product ID: {}",
event["type"],
product_data["product_id"],
)
create_or_update_product(session, product_data)
elif product_is_deleted(event):
delete_product_by_id(session, event["data"]["product_id"])


def _handle_interaction_event(event: dict) -> None:
with get_session() as session:
record_interaction(session, event)
product = get_product_by_id(session, event["data"]["product_id"])
if not product:
try:
product_data = fetch_product_by_id(event["data"]["product_id"])
create_product(session, product_data)
session.commit()
except requests.RequestException as exc:
logger.error(
"Failed to fetch product {} for interaction event {}: {}",
event["data"]["product_id"],
event["type"],
exc,
)
session.rollback()

if __name__ == "__main__":
start_kafka_consumer()
94 changes: 0 additions & 94 deletions recommender/app/entry/main.py

This file was deleted.

64 changes: 0 additions & 64 deletions recommender/app/entry/sync.py

This file was deleted.

Loading
Loading