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
130 changes: 89 additions & 41 deletions src/main/java/com/techevents/consumer/EventConsumer.java
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,12 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import com.techevents.model.Event;
import com.techevents.repository.EventRepository;
import com.techevents.service.EventIngestionService;
import com.techevents.config.KafkaTopics;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import org.springframework.dao.DataAccessException;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
Expand All @@ -22,51 +21,100 @@ public class EventConsumer {

private static final Logger log = LoggerFactory.getLogger(EventConsumer.class);

private final ObjectMapper objectMapper;
private final EventRepository eventRepository;
private final ObjectMapper mapper;
private final EventIngestionService eventIngestionService;

public EventConsumer(ObjectMapper objectMapper, EventRepository eventRepository) {
this.objectMapper = objectMapper.registerModule(new JavaTimeModule());
this.eventRepository = eventRepository;
public EventConsumer(ObjectMapper mapper,
EventIngestionService ingestion) {
this.mapper = mapper.registerModule(new JavaTimeModule());
this.eventIngestionService = ingestion;
}

@KafkaListener(
topics = KafkaTopics.EVENTS_INGEST,
groupId = "event-group",
containerFactory = "kafkaListenerContainerFactory"
)
public void listen(
@Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
@Payload String message
) {
log.debug("Consuming message from {}-{}-{}", topic, partition, offset);

// 1) DESERIALIZE
Event event;
@KafkaListener(topics = KafkaTopics.EVENTS_INGEST, groupId = "event-group")
public void listen(@Header(KafkaHeaders.RECEIVED_TOPIC) String topic, @Payload String payload) {
try {
event = objectMapper.readValue(message, Event.class);
log.info("Parsed Event[id={}, title={}]", event.getId(), event.getTitle());
} catch (JsonProcessingException e) {
log.error("JSON parse error at {}-{}-{}; payload={}", topic, partition, offset, message, e);
return;
}
Event e = mapper.readValue(payload, Event.class);

// 2) VALIDATE
if (event.getTitle() == null || event.getEventDate() == null) {
log.warn("Validation failed for {}-{}-{}; payload={}", topic, partition, offset, message);
return;
}
// Optional validation (can move to service)
if (e.getTitle() == null || e.getEventDate() == null) {
log.warn("Invalid payload received from topic {}: {}", topic, payload);
return;
}

// 3) PERSIST (let DB errors bubble up)
try {
eventRepository.save(event);
log.info("Saved Event[id={}, title={}]", event.getId(), event.getTitle());
} catch (DataAccessException dbEx) {
log.error("DB error saving Event[id={}, title={}] at {}-{}-{}",
event.getId(), event.getTitle(), topic, partition, offset, dbEx);
throw dbEx;
eventIngestionService.ingest(e);

} catch (JsonProcessingException ex) {
log.error("Failed to deserialize payload from topic {}: {}", topic, payload, ex);
} catch (Exception ex) {
log.error("Unexpected error while processing payload from topic {}: {}", topic, payload, ex);
}
}

}

// @Component
// public class EventConsumer {

// private static final Logger log =
// LoggerFactory.getLogger(EventConsumer.class);

// private final ObjectMapper objectMapper;
// private final EventRepository eventRepository;

// private final EventIngestionService eventIngestionService;

// public EventConsumer(ObjectMapper objectMapper, EventRepository
// eventRepository, EventIngestionService eventIngestionService) {
// this.objectMapper = objectMapper.registerModule(new JavaTimeModule());
// this.eventRepository = eventRepository;
// this.eventIngestionService = eventIngestionService;
// }

// @KafkaListener(
// topics = KafkaTopics.EVENTS_INGEST,
// groupId = "event-group",
// containerFactory = "kafkaListenerContainerFactory"
// )
// public void listen(
// @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
// @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
// @Header(KafkaHeaders.OFFSET) long offset,
// @Payload String message
// ) {
// log.debug("Consuming message from {}-{}-{}", topic, partition, offset);

// // 1) DESERIALIZE
// Event event;
// try {
// event = objectMapper.readValue(message, Event.class);
// log.info("Parsed Event[id={}, title={}]", event.getId(), event.getTitle());
// } catch (JsonProcessingException e) {
// log.error("JSON parse error at {}-{}-{}; payload={}", topic, partition,
// offset, message, e);
// return;
// }

// // 2) VALIDATE
// if (event.getTitle() == null || event.getEventDate() == null) {
// log.warn("Validation failed for {}-{}-{}; payload={}", topic, partition,
// offset, message);
// return;
// }

// // 3) PERSIST (let DB errors bubble up)
// try {
// eventRepository.save(event);
// log.info("Saved Event[id={}, title={}]", event.getId(), event.getTitle());

// eventIngestionService.processEvent(event);
// log.info("Processed Event[id={}, title={}]", event.getId(),
// event.getTitle());

// } catch (DataAccessException dbEx) {
// log.error("DB error saving Event[id={}, title={}] at {}-{}-{}",
// event.getId(), event.getTitle(), topic, partition, offset, dbEx);
// throw dbEx;
// }

// }
// }
35 changes: 35 additions & 0 deletions src/main/java/com/techevents/model/NotificationMessage.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
package com.techevents.model;

import com.fasterxml.jackson.annotation.JsonInclude;

@JsonInclude(JsonInclude.Include.NON_NULL)
public class NotificationMessage {
private String email;
private String eventTitle;
private String eventDate;
private String city;

public NotificationMessage(String email, Event event) {
this.email = email;
this.eventTitle = event.getTitle();
this.eventDate = event.getEventDate().toString();
this.city = event.getCity();
}

// Getters (and setters if needed)
public String getEmail() {
return email;
}

public String getEventTitle() {
return eventTitle;
}

public String getEventDate() {
return eventDate;
}

public String getCity() {
return city;
}
}
45 changes: 45 additions & 0 deletions src/main/java/com/techevents/service/EventIngestionService.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package com.techevents.service;

import com.techevents.model.Event;
import com.techevents.repository.EventRepository;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;

@Service
public class EventIngestionService {

private static final Logger log = LoggerFactory.getLogger(EventIngestionService.class);

private final EventRepository eventRepository;
private final RabbitTemplate rabbit;

@Value("${app.rabbitmq.exchange}")
private String exchange;

@Value("${app.rabbitmq.routing-key}")
private String routingKey;

public EventIngestionService(EventRepository eventRepository, RabbitTemplate rabbitTemplate) {
this.eventRepository = eventRepository;
this.rabbit = rabbitTemplate;
}

public void ingest(Event event) {
// Validate
if (event.getTitle() == null || event.getEventDate() == null) {
throw new IllegalArgumentException("Event must have a title and an event date.");
}

// 1) Persist to database
Event saved = eventRepository.save(event);
log.info("Saved Event[id={}, title={}]", saved.getId(), saved.getTitle());

// 2) Publish to RabbitMQ
rabbit.convertAndSend(exchange, routingKey, saved);
log.info("Published Event[id={}] to RabbitMQ [exchange={}, routingKey={}]",
saved.getId(), exchange, routingKey);
}
}
44 changes: 44 additions & 0 deletions src/main/java/com/techevents/service/NotificationService.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
package com.techevents.service;

import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.techevents.model.NotificationMessage;
import com.techevents.model.Subscriber;
import com.techevents.model.Event;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;

@Service
public class NotificationService {

private static final Logger log = LoggerFactory.getLogger(NotificationService.class);

private final RabbitTemplate rabbitTemplate;
private final ObjectMapper objectMapper;

@Value("${app.rabbitmq.notification-exchange}")
private String exchange;

@Value("${app.rabbitmq.notification-routing-key}")
private String routingKey;

public NotificationService(RabbitTemplate rabbitTemplate, ObjectMapper objectMapper) {
this.rabbitTemplate = rabbitTemplate;
this.objectMapper = objectMapper;
}

public void sendNotification(Subscriber subscriber, Event event) {
NotificationMessage message = new NotificationMessage(subscriber.getEmail(), event);

try {
String jsonPayload = objectMapper.writeValueAsString(message);
rabbitTemplate.convertAndSend(exchange, routingKey, jsonPayload);
log.info("Published notification for email={} to RabbitMQ", subscriber.getEmail());
} catch (JsonProcessingException e) {
log.error("Failed to serialize NotificationMessage for subscriber={}", subscriber.getEmail(), e);
}
}
}
9 changes: 6 additions & 3 deletions src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,12 @@ spring:

app:
rabbitmq:
queue: test.queue
exchange: test.exchange
routing-key: test.key
queue: events.saved.queue
exchange: events.saved.exchange
routing-key: events.saved.key
notification-exchange: notifications.exchange
notification-routing-key: notifications.event.created


server:
port: 8080
Loading