PaymentEventListener.java
Bridget/payment-gateway-microservices/services/notification-service/src/main/java/com/checkout/paymentgw/notify/consumer/PaymentEventListener.java
package com.checkout.paymentgw.notify.consumer;
import com.checkout.paymentgw.events.PaymentEvent;
import com.checkout.paymentgw.events.Topics;
import com.checkout.paymentgw.notify.domain.NotificationLog;
import com.checkout.paymentgw.notify.domain.NotificationLogRepository;
import com.checkout.paymentgw.notify.notifier.Notifier;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.time.Instant;
/**
* Idempotent consumer. The {@code event_id} is the primary key of {@code notification_log};
* a duplicate insert raises a unique-violation that we treat as "already processed".
*/
@Component
public class PaymentEventListener {
private static final Logger LOG = LoggerFactory.getLogger(PaymentEventListener.class);
private final NotificationLogRepository repository;
private final Notifier notifier;
private final ObjectMapper mapper;
public PaymentEventListener(NotificationLogRepository repository,
Notifier notifier,
ObjectMapper mapper) {
this.repository = repository;
this.notifier = notifier;
this.mapper = mapper;
}
@KafkaListener(topics = Topics.PAYMENTS_EVENTS, groupId = "notification-service")
@Transactional
public void onMessage(String payload) {
PaymentEvent ev;
try {
ev = mapper.readValue(payload, PaymentEvent.class);
} catch (Exception e) {
LOG.error("cannot deserialise event payload, skipping; payload size={}", payload.length(), e);
return;
}
if (ev.eventId() == null) {
LOG.warn("event without eventId, skipping");
return;
}
try {
repository.save(NotificationLog.builder()
.eventId(ev.eventId())
.paymentId(ev.paymentId())
.eventType(ev.eventType())
.receivedAt(Instant.now())
.build());
} catch (DataIntegrityViolationException dup) {
LOG.info("duplicate event {} — already processed", ev.eventId());
return;
}
notifier.deliver(ev);
}
}
Articles liés
PaymentGatewayApplication.java
PaymentGatewayApplication.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/PaymentGatewayApplication.java).
Lire l'article →ApplicationConfiguration.java
ApplicationConfiguration.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/configuration/ApplicationConfiguration.java).
Lire l'article →RequestResponseLogger.java
RequestResponseLogger.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/configuration/RequestResponseLogger.java).
Lire l'article →PaymentGatewayController.java
PaymentGatewayController.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/controller/PaymentGatewayController.java).
Lire l'article →PaymentStatus.java
PaymentStatus.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/enums/PaymentStatus.java).
Lire l'article →BankErrorException.java
BankErrorException.java — java source code from the Bridget learning materials (Bridget/payment-gateway-challenge-java/src/main/java/com/checkout/payment/gateway/exception/BankErrorException.java).
Lire l'article →