S SmartDocs
系列: Bridget java 58 行 · 更新于 2026-05-08

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
            }
        }
    }
}

相关文章

Bridget java 更新于 2026-03-30

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).

阅读文章 →
Bridget java 更新于 2026-03-20

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).

阅读文章 →
Bridget java 更新于 2026-03-30

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).

阅读文章 →
Bridget java 更新于 2026-04-14

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).

阅读文章 →
Bridget java 更新于 2026-03-20

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).

阅读文章 →
Bridget java 更新于 2026-03-29

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).

阅读文章 →