Рефакторинг SomeServiceImpl для асинхронности и корректности

7. Исправить SomeServiceImpl и вынести долгий вызов налоговой в асинхронную обработку

Условие задачи:
Дан высоконагруженный Spring-микросервис, запущенный в Kubernetes в 8 репликах.

Сервис выполняет две независимые задачи:

  1. processJobRunning() раз в час получает из PostgreSQL необработанные возвратные чеки, отправляет событие в аналитику и помечает чек обработанным.

  2. sendReceipt() вызывается по HTTP другим микросервисом и отправляет чек во внешнюю налоговую через Feign-клиент. Налоговая может отвечать до 60 секунд, поэтому HTTP-поток не должен ждать её ответа.

Необходимо провести code review и исправить проблемы с AOP, транзакциями, кэшированием, конкурентным выполнением scheduled-задачи и долгим внешним вызовом.

Код:

@Service
@AllArgsConstructor
public class SomeServiceImpl implements SomeService {

    private final TaxClient taxClient;
    private final ReceiptDao receiptDao;
    private final ReceiptTypeDao receiptTypeDao;

    private final ReceiptMapper mapper;
    private final ReceiptTypeMapper typeMapper;
    private final RefundReceiptMapper refundReceiptMapper;

    private final SomeService self;

    @Scheduled(cron = "${cron.expression}")
    public void processJobRunning() {
        List<Receipt> receipts = getRefundReceipts();

        for (Receipt receipt : receipts) {
            receipt.setProcessed(true);

            eventListener.sendEventToAnalytics(
                    new SomeEvent(receipt.getId())
            );

            receiptDao.save(
                    mapReceipt(mapper.map(receipt))
            );
        }
    }

    @Transactional(readOnly = true)
    public List<Receipt> getRefundReceipts() {
        return receiptDao.findAllBySourceAndProcessedFalse(
                ReceiptSource.REFUND
        );
    }

    private Receipt mapReceipt(ReceiptDto receiptDto) {
        return mapper.map(receiptDto);
    }

    public void sendReceipt(ReceiptDto receiptDto) {
        ReceiptSource source =
                self.getReceiptSource(receiptDto);

        taxClient.sendReceipt(receiptDto, source);
    }

    @Transactional
    @Cacheable(
            value = "receipt_type",
            cacheManager = "redisCacheManager"
    )
    public ReceiptSource getReceiptSource(
            ReceiptDto receiptDto
    ) {
        Receipt receipt =
                receiptDao.findById(receiptDto.getId());

        return receipt.getSource();
    }
}

Спойлеры к решению

Подсказки
💡 findById() возвращает Optional<Receipt>.
💡 Вызов метода с @Transactional или @Cacheable из того же объекта обходит Spring AOP-прокси.
💡 Инъекция собственного интерфейса через конструктор — плохой способ обходить self-invocation и может привести к циклической зависимости.
💡 @Scheduled выполнится на каждой из восьми реплик — одного @Scheduled недостаточно.
💡 @Async освобождает HTTP-поток, но Feign-вызов всё равно блокирует отдельный worker-thread.
💡 Для высоконагруженной системы надёжнее очередь сообщений: Kafka/RabbitMQ + отдельный consumer.

Решение

Здесь несколько независимых проблем.

1. findById() используется неверно

Spring Data возвращает:

Optional<Receipt>

поэтому запись нужно получать явно:

receiptDao.findById(receiptId)
        .orElseThrow(...);

2. @Cacheable и @Transactional не работают при self-invocation

Вызов:

this.getReceiptSource(...)

не проходит через Spring AOP proxy.

Инъекция:

private final SomeService self;

тоже нежелательна: она создаёт зависимость бина от самого себя и может привести к циклической зависимости.

Лучше вынести получение источника в отдельный Spring-бин:

@Service
@RequiredArgsConstructor
public class ReceiptSourceService {

    private final ReceiptDao receiptDao;

    @Transactional(readOnly = true)
    @Cacheable(
            cacheNames = "receipt_type",
            key = "#receiptId",
            cacheManager = "redisCacheManager"
    )
    public ReceiptSource getReceiptSource(Long receiptId) {
        return receiptDao.findById(receiptId)
                .orElseThrow(() ->
                        new EntityNotFoundException(
                                "Receipt not found: " + receiptId
                        )
                )
                .getSource();
    }
}

Теперь вызов идёт между двумя Spring-бинами, поэтому и @Cacheable, и @Transactional применяются через proxy.

Также явно задан ключ:

key = "#receiptId"

В противном случае Spring использовал бы в качестве ключа весь ReceiptDto, что для такого кэша не требуется.


3. Долгий Feign-вызов нельзя выполнять в HTTP-потоке

Минимальный вариант — вынести его в отдельный бин с @Async:

@Service
@RequiredArgsConstructor
@Slf4j
public class TaxReceiptSender {

    private final TaxClient taxClient;

    @Async("taxExecutor")
    public void send(
            ReceiptDto receipt,
            ReceiptSource source
    ) {
        try {
            taxClient.sendReceipt(receipt, source);
        } catch (Exception e) {
            log.error(
                    "Failed to send receipt {} to tax service",
                    receipt.getId(),
                    e
            );

            throw e;
        }
    }
}

Тогда основной сервис выглядит так:

@Service
@RequiredArgsConstructor
@Slf4j
public class SomeServiceImpl implements SomeService {

    private final ReceiptDao receiptDao;
    private final SomeEventListener eventListener;
    private final ReceiptSourceService receiptSourceService;
    private final TaxReceiptSender taxReceiptSender;

    @Override
    public void sendReceipt(ReceiptDto receiptDto) {
        Objects.requireNonNull(receiptDto, "receiptDto");
        Objects.requireNonNull(receiptDto.getId(), "receiptDto.id");

        ReceiptSource source =
                receiptSourceService.getReceiptSource(
                        receiptDto.getId()
                );

        taxReceiptSender.send(receiptDto, source);
    }
}

sendReceipt() теперь не ждёт до 60 секунд ответа налоговой. Он только ставит задачу на выполнение в taxExecutor.

Для @Async необходимо включить асинхронность:

@Configuration
@EnableAsync
public class AsyncConfig {

    @Bean("taxExecutor")
    public Executor taxExecutor() {
        ThreadPoolTaskExecutor executor =
                new ThreadPoolTaskExecutor();

        executor.setCorePoolSize(16);
        executor.setMaxPoolSize(32);
        executor.setQueueCapacity(500);
        executor.setThreadNamePrefix("tax-");

        executor.initialize();

        return executor;
    }
}

Размеры пула и очереди здесь примерные — в реальном сервисе их необходимо выбирать исходя из нагрузки и возможностей налоговой.

Важно: @Async не делает Feign неблокирующим.

Получается:

HTTP thread
    ├── submit task
    └── быстро освобождается

taxExecutor thread
    └── может ждать Feign до 60 секунд

Поэтому использовать общий executor Spring нельзя — для такого долгого внешнего вызова нужен отдельный ограниченный пул.


4. @Scheduled в восьми pod’ах выполнится восемь раз

Если запущено:

pod-1
pod-2
...
pod-8

то каждый pod независимо вызовет:

@Scheduled
public void processJobRunning()

В результате один и тот же чек могут одновременно обработать несколько реплик.

Для scheduled-задачи нужен распределённый lock, например через ShedLock с общим хранилищем в PostgreSQL или Redis:

@Scheduled(cron = "${cron.expression}")
@SchedulerLock(
        name = "refund-receipts-job",
        lockAtMostFor = "PT2H"
)
public void processJobRunning() {
    List<Receipt> receipts =
            receiptDao.findAllBySourceAndProcessedFalse(
                    ReceiptSource.REFUND
            );

    for (Receipt receipt : receipts) {
        try {
            eventListener.sendEventToAnalytics(
                    new SomeEvent(receipt.getId())
            );

            receipt.setProcessed(true);
            receiptDao.save(receipt);

        } catch (Exception e) {
            log.error(
                    "Failed to process receipt {}",
                    receipt.getId(),
                    e
            );
        }
    }
}

Теперь одновременно job выполняет только одна реплика.

Если требуется, наоборот, распределять обработку между всеми pod’ами, вместо глобального lock лучше применять захват записей в PostgreSQL через FOR UPDATE SKIP LOCKED или отдельный статус PROCESSING.


5. processed лучше устанавливать после успешной отправки события

Последовательность:

eventListener.sendEventToAnalytics(...);

receipt.setProcessed(true);
receiptDao.save(receipt);

лучше отражает бизнес-смысл: чек считается обработанным только после успешной отправки события.

Если отправка события прошла, а сохранение processed = true упало, событие может отправиться повторно на следующем запуске.

По условию событие идемпотентно, поэтому такая семантика at-least-once допустима.


6. Лишнее преобразование Receipt → DTO → Receipt

Конструкция:

receiptDao.save(
        mapReceipt(mapper.map(receipt))
);

не нужна.

Уже имеется объект:

Receipt receipt

поэтому достаточно:

receiptDao.save(receipt);

Преобразование:

Receipt
ReceiptDto
Receipt

только создаёт лишние объекты и потенциально может потерять часть данных сущности.


7. getRefundReceipts() с @Transactional тоже страдал от self-invocation

Вызов:

processJobRunning()
    
getRefundReceipts()

происходил внутри того же объекта, поэтому @Transactional на getRefundReceipts() через Spring AOP не применялся.

В данном случае отдельный метод вообще необязателен: Spring Data repository самостоятельно выполняет запрос в необходимом repository transaction context.


8. Cron-выражение

Правильная запись:

@Scheduled(cron = "${cron.expression}")

Для Spring cron обычно содержит шесть полей.

Например, запуск каждый час в 30:30:

30 30 * * * *

а не пятичастное:

30 30 * * *

Итоговый основной сервис:

@Service
@RequiredArgsConstructor
@Slf4j
public class SomeServiceImpl implements SomeService {

    private final ReceiptDao receiptDao;
    private final SomeEventListener eventListener;
    private final ReceiptSourceService receiptSourceService;
    private final TaxReceiptSender taxReceiptSender;

    @Scheduled(cron = "${cron.expression}")
    @SchedulerLock(
            name = "refund-receipts-job",
            lockAtMostFor = "PT2H"
    )
    public void processJobRunning() {
        List<Receipt> receipts =
                receiptDao.findAllBySourceAndProcessedFalse(
                        ReceiptSource.REFUND
                );

        for (Receipt receipt : receipts) {
            try {
                eventListener.sendEventToAnalytics(
                        new SomeEvent(receipt.getId())
                );

                receipt.setProcessed(true);
                receiptDao.save(receipt);

            } catch (Exception e) {
                log.error(
                        "Failed to process receipt {}",
                        receipt.getId(),
                        e
                );
            }
        }
    }

    @Override
    public void sendReceipt(ReceiptDto receiptDto) {
        Objects.requireNonNull(receiptDto, "receiptDto");
        Objects.requireNonNull(receiptDto.getId(), "receiptDto.id");

        ReceiptSource source =
                receiptSourceService.getReceiptSource(
                        receiptDto.getId()
                );

        taxReceiptSender.send(receiptDto, source);
    }
}

Для высоконагруженного production-сервиса вместо @Async предпочтительнее очередь сообщений:

HTTP request
sendReceipt()
Kafka / RabbitMQ
ответ вызывающему сервису

consumer
TaxClient
Налоговая

То есть sendReceipt() только публикует команду:

public void sendReceipt(ReceiptDto receiptDto) {
    ReceiptSource source =
            receiptSourceService.getReceiptSource(
                    receiptDto.getId()
            );

    kafkaTemplate.send(
            "tax-receipts",
            new TaxReceiptCommand(receiptDto, source)
    );
}

А отдельный consumer выполняет медленный вызов:

@KafkaListener(topics = "tax-receipts")
public void process(TaxReceiptCommand command) {
    taxClient.sendReceipt(
            command.receipt(),
            command.source()
    );
}

Для данного сценария это надёжнее @Async, потому что:

  • запрос не ждёт налоговую;

  • задача не теряется при рестарте pod’а;

  • можно организовать retry;

  • можно использовать DLQ;

  • Kafka обеспечивает backpressure;

  • восемь pod’ов могут параллельно работать как consumer group.

Поэтому @Async — нормальное минимальное исправление, а очередь сообщений — предпочтительный вариант для описанного высоконагруженного микросервиса.