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-ограничением, чтобы повторная обработка не создавала вторую запись.