OutboxPublisher.java
Bridget/payment-gateway-microservices/services/payment-service/src/main/java/com/checkout/paymentgw/payment/outbox/OutboxPublisher.java
package com.checkout.paymentgw.payment.outbox;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.data.domain.PageRequest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import java.time.Instant;
import java.util.List;
/**
* Drains unpublished outbox rows to Kafka. At-least-once: a row stays "unpublished"
* until the broker has acknowledged the produce.
*
* <p>Ordering: rows are read in {@code created_at ASC} and published serially per batch.
* For per-aggregate ordering we partition by {@link OutboxMessage#getKeyHint()}.
*/
@Component
public class OutboxPublisher {
private static final Logger LOG = LoggerFactory.getLogger(OutboxPublisher.class);
private final OutboxRepository outboxRepository;
private final KafkaTemplate<String, String> kafkaTemplate;
private final int batchSize;
public OutboxPublisher(OutboxRepository outboxRepository,
KafkaTemplate<String, String> kafkaTemplate) {
this.outboxRepository = outboxRepository;
this.kafkaTemplate = kafkaTemplate;
this.batchSize = 50;
}
@Scheduled(fixedDelayString = "${payment.outbox.poll-interval:1000}")
@Transactional
public void drain() {
List<OutboxMessage> batch = outboxRepository.findUnpublished(PageRequest.of(0, batchSize));
if (batch.isEmpty()) {
return;
}
LOG.debug("draining {} outbox messages", batch.size());
for (OutboxMessage msg : batch) {
try {
kafkaTemplate.send(msg.getTopic(), msg.getKeyHint(), msg.getPayload()).get();
msg.setPublishedAt(Instant.now());
outboxRepository.save(msg);
} catch (Exception e) {
// Leave row unpublished — we'll retry next tick. Bounded retry budget would
// live here in a production system (with a dead-letter table).
LOG.warn("publish failed for outbox id={}: {}", msg.getId(), e.getMessage());
break; // back off — subsequent rows likely fail too
}
}
}
}
Artículos relacionados
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).
Leer artículo →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).
Leer artículo →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).
Leer artículo →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).
Leer artículo →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).
Leer artículo →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).
Leer artículo →