Сделать ревью потокобезопасного сервиса сохранения заданий и работ

20. Сделать ревью потокобезопасного сервиса сохранения заданий и работ

Условие задачи:
Дан Spring-сервис, который должен сохранять сущности двух типов — Task и Work. Сервис может вызываться одновременно из REST-контроллера и Kafka listener.

Нужно провести code review, проверить корректность многопоточной работы и выполнить рефакторинг так, чтобы сервис безопасно обрабатывал параллельные вызовы и корректно сохранял сущность нужного типа.

Код:

@Component
public class MyShinyService {

    @Autowired
    private final TaskRepository taskRepository;

    @Autowired
    private final WorkRepository workRepository;

    // for concurrency safety
    private final Task task;

    @Transactional
    public synchronized accumulateTask(Task task) {
        this.task = task;
    }

    @Transactional
    public synchronized void saveTask() {
        System.out.println(""saving process has started"");

        String task = ""task"";

        if (this.task.type == task) {
            this.taskRepository.save(
                    new TaskEntity(
                            UUID.randomUUID(),
                            this.task.id.toString(),
                            this.task.name
                    )
            );
        } else {
            this.workRepository.save(
                    new WorkEntity(
                            UUID.randomUUID(),
                            this.task.id.toString(),
                            this.task.name
                    )
            );
        }

        System.out.println(""saved"");
    }
}

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

Подсказки
💡 Проверь, компилируются ли final-поля и accumulateTask().
💡 Spring-сервис обычно singleton, поэтому хранить текущую задачу в поле опасно.
💡 Два synchronized-метода не делают цепочку из двух вызовов атомарной.
💡 Для REST и Kafka стоит учитывать повторную обработку одного задания.

Решение

Проблемы:

  • accumulateTask() не имеет возвращаемого типа;

  • taskRepository, workRepository и task объявлены final, но не инициализируются конструктором;

  • this.task = task невозможно для final-поля;

  • сервис хранит изменяемое состояние в singleton-бине;

  • synchronized защищает каждый метод отдельно, но не всю последовательность accumulateTask()saveTask();

  • saveTask() может быть вызван до установки task;

  • synchronized не защищает между несколькими экземплярами приложения;

  • строки сравниваются через ==;

  • @Transactional на accumulateTask() не нужен — там нет работы с БД;

  • field injection лучше заменить constructor injection;

  • строковый тип лучше заменить на enum;

  • System.out.println() лучше заменить логгером;

  • при Kafka redelivery или REST retry возможны дубликаты.

Главная проблема — состояние внутри сервиса.

Например:

Thread 1: accumulateTask(task1)
Thread 2: accumulateTask(task2)
Thread 1: saveTask()

К моменту saveTask() в поле уже находится task2. Поэтому безопаснее вообще не хранить текущую задачу в сервисе, а передавать её непосредственно в метод сохранения.

@Service
@RequiredArgsConstructor
public class MyShinyService {

    private final TaskRepository taskRepository;
    private final WorkRepository workRepository;

    @Transactional
    public UUID save(Task task) {
        validate(task);

        UUID entityId = UUID.randomUUID();

        switch (task.getType()) {
            case TASK -> taskRepository.save(
                    new TaskEntity(
                            entityId,
                            task.getId().toString(),
                            task.getName()
                    )
            );

            case WORK -> workRepository.save(
                    new WorkEntity(
                            entityId,
                            task.getId().toString(),
                            task.getName()
                    )
            );
        }

        return entityId;
    }

    private void validate(Task task) {
        if (task == null
                || task.getId() == null
                || task.getType() == null
                || task.getName() == null
                || task.getName().isBlank()) {
            throw new IllegalArgumentException(
                    "Invalid task"
            );
        }
    }
}

Тип лучше задать явно:

public enum TaskType {
    TASK,
    WORK
}

Теперь сервис stateless: каждый вызов работает только со своим объектом, поэтому synchronized не нужен.

Если метод вызывается из Kafka или REST с возможными retry, дополнительно нужна идемпотентность. Например, исходный task.id можно хранить с UNIQUE-ограничением, чтобы повторная обработка не создавала вторую запись.