Event-driven: TaskCreated → Audit log в БД
Цель урока
Сделать полноценный CRUD для задач с правильными HTTP-статусами (201/204/404), обработкой ошибок, и подготовить Task к подключению JPA.
Теория · для собеса
1 Зачем audit log в БД ▸
- Compliance (GDPR, SOX, HIPAA) — нужно знать, кто менял данные
- Debugging — «когда статус задачи стал DONE? Кто?»
- Time-travel — можно восстановить состояние на любой момент
- Forensics — если юзер жалуется «я ничего не удалял», audit log покажет
2 Структура audit event ▸
CREATE TABLE task_events (
id BIGSERIAL PRIMARY KEY,
task_id BIGINT NOT NULL,
event_type VARCHAR(20) NOT NULL, -- CREATED, UPDATED, DELETED
event_data JSONB, -- полный snapshot или diff
actor_username VARCHAR(50) NOT NULL,
occurred_at TIMESTAMP NOT NULL,
received_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);
CREATE INDEX idx_task_events_task_id ON task_events(task_id);
CREATE INDEX idx_task_events_occurred_at ON task_events(occurred_at);
⚠️ JSONB (PostgreSQL) — бинарный JSON, можно индексировать. Не JSON! (текстовый, медленнее).
3 Outbox pattern (продвинутая тема) ▸
Проблема: если БД-запись прошла, а Kafka-отправка упала — событие потеряно.
Решение: пишем в БД в той же транзакции + отдельный процесс публикует в Kafka.
Out of scope для собеса, но упомяни на интервью как «знаю, но не реализовал».
---
Практика: улучшаем CRUD
-
1Flyway-миграция
V4__add_task_events.sqlCREATE TABLE task_events ( id BIGSERIAL PRIMARY KEY, task_id BIGINT NOT NULL, event_type VARCHAR(20) NOT NULL, actor_username VARCHAR(50) NOT NULL, event_data JSONB, occurred_at TIMESTAMP NOT NULL, received_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX idx_task_events_task_id ON task_events(task_id); CREATE INDEX idx_task_events_actor ON task_events(actor_username); CREATE INDEX idx_task_events_occurred_at ON task_events(occurred_at DESC); -
2
model/TaskEventLog.java(Entity)package com.taskflow.model; import com.fasterxml.jackson.databind.JsonNode; import io.hypersistence.utils.hibernate.type.json.JsonBinaryType; import jakarta.persistence.*; import org.hibernate.annotations.Type; import java.time.Instant; @Entity @Table(name = "task_events") public class TaskEventLog { @Id @GeneratedValue(strategy = GenerationType.IDENTITY) private Long id; @Column(name = "task_id", nullable = false) private Long taskId; @Column(name = "event_type", nullable = false, length = 20) private String eventType; @Column(name = "actor_username", nullable = false, length = 50) private String actorUsername; @Type(JsonBinaryType.class) @Column(name = "event_data", columnDefinition = "jsonb") private JsonNode eventData; @Column(name = "occurred_at", nullable = false) private Instant occurredAt; @Column(name = "received_at", nullable = false, updatable = false) private Instant receivedAt; public TaskEventLog() {} @PrePersist void onCreate() { this.receivedAt = Instant.now(); } // Getters / Setters public Long getId() { return id; } public void setId(Long id) { this.id = id; } public Long getTaskId() { return taskId; } public void setTaskId(Long taskId) { this.taskId = taskId; } public String getEventType() { return eventType; } public void setEventType(String eventType) { this.eventType = eventType; } public String getActorUsername() { return actorUsername; } public void setActorUsername(String actorUsername) { this.actorUsername = actorUsername; } public JsonNode getEventData() { return eventData; } public void setEventData(JsonNode eventData) { this.eventData = eventData; } public Instant getOccurredAt() { return occurredAt; } public void setOccurredAt(Instant occurredAt) { this.occurredAt = occurredAt; } public Instant getReceivedAt() { return receivedAt; } }Зависимость для
JsonBinaryType(Hibernate Types):dependencies { implementation("io.hypersistence:hypersistence-utils-hibernate-63:3.8.3") } -
3
repository/TaskEventLogRepository.javapackage com.taskflow.repository; import com.taskflow.model.TaskEventLog; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Repository; import java.util.List; @Repository public interface TaskEventLogRepository extends JpaRepository{ List findByTaskIdOrderByOccurredAtDesc(Long taskId); List findByActorUsernameOrderByOccurredAtDesc(String username); } -
4Расширяем
event/TaskCreatedEvent— добавляем snapshotpublic record TaskCreatedEvent( Long taskId, String title, String description, boolean done, String ownerUsername, Instant occurredAt ) implements TaskEvent { @Override public String eventType() { return "TaskCreated"; } } -
5
consumer/TaskEventConsumer.java— пишем в БДpackage com.taskflow.consumer; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.taskflow.event.TaskCreatedEvent; import com.taskflow.event.TaskDeletedEvent; import com.taskflow.event.TaskEvent; import com.taskflow.event.TaskUpdatedEvent; import com.taskflow.model.TaskEventLog; import com.taskflow.repository.TaskEventLogRepository; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import org.springframework.transaction.annotation.Transactional; @Component public class TaskEventConsumer { private static final Logger log = LoggerFactory.getLogger(TaskEventConsumer.class); private final TaskEventLogRepository repo; private final ObjectMapper mapper; public TaskEventConsumer(TaskEventLogRepository repo, ObjectMapper mapper) { this.repo = repo; this.mapper = mapper; } @KafkaListener(topics = "task-events", groupId = "taskflow-audit-writer") @Transactional public void onTaskEvent(TaskEvent event) { try { TaskEventLog logEntry = new TaskEventLog(); logEntry.setTaskId(event.taskId()); logEntry.setEventType(event.eventType()); logEntry.setOccurredAt(event.occurredAt()); logEntry.setEventData(mapper.valueToTree(event)); // Извлекаем actor String actor = switch (event) { case TaskCreatedEvent e -> e.ownerUsername(); case TaskUpdatedEvent e -> e.ownerUsername() != null ? e.ownerUsername() : "system"; case TaskDeletedEvent e -> "system"; }; logEntry.setActorUsername(actor); repo.save(logEntry); log.info("Audit log written: task={}, event={}, actor={}", event.taskId(), event.eventType(), actor); } catch (Exception e) { log.error("Failed to write audit log for event: {}", event, e); // В проде: отправить в DLQ (Dead Letter Queue) throw e; // Kafka ретраит } } } -
6
service/AuditService.java— отдаём audit logpackage com.taskflow.service; import com.taskflow.dto.AuditEventResponse; import com.taskflow.model.TaskEventLog; import com.taskflow.repository.TaskEventLogRepository; import org.springframework.stereotype.Service; import java.util.List; @Service public class AuditService { private final TaskEventLogRepository repo; public AuditService(TaskEventLogRepository repo) { this.repo = repo; } public ListgetTaskHistory(Long taskId) { return repo.findByTaskIdOrderByOccurredAtDesc(taskId) .stream() .map(AuditEventResponse::from) .toList(); } } -
7DTO
public record AuditEventResponse( Long id, Long taskId, String eventType, String actorUsername, JsonNode eventData, Instant occurredAt ) { public static AuditEventResponse from(TaskEventLog log) { return new AuditEventResponse(log.getId(), log.getTaskId(), log.getEventType(), log.getActorUsername(), log.getEventData(), log.getOccurredAt()); } } -
8Controller
@RestController @RequestMapping("/api/tasks/{taskId}/audit") public class AuditController { private final AuditService auditService; public AuditController(AuditService auditService) { this.auditService = auditService; } @GetMapping @PreAuthorize("hasRole('ADMIN') or @taskSecurity.isOwner(#taskId, authentication.name)") public ListgetHistory(@PathVariable Long taskId) { return auditService.getTaskHistory(taskId); } } -
9Тест
# Создаём задачу $ curl -X POST http://localhost:8080/api/tasks \ -H "Authorization: Bearer $TOKEN" \ -H "Content-Type: application/json" \ -d '{"title":"Audit test"}' # Видим audit log $ curl -H "Authorization: Bearer $TOKEN" http://localhost:8080/api/tasks/1/audit [ { "id": 1, "taskId": 1, "eventType": "TaskCreated", "actorUsername": "alice", "eventData": {"taskId":1,"title":"Audit test","description":null,"done":false,"ownerUsername":"alice"}, "occurredAt": "2026-06-29T16:55:00Z" } ] # Проверяем в БД $ docker exec -it taskflow-postgres psql -U taskflow -d taskflow_dev -c \ "SELECT event_type, actor_username, occurred_at FROM task_events ORDER BY occurred_at DESC LIMIT 5;"---
Зачем это на собесе
После урока ты должен уметь ответить на:
JSONB индексируется в PostgreSQL — можно искать по содержимому события (быстро).task-events-dlq для ручной обработки.5 вопросов на углубление
Раскрой вопрос и нажми «🤔 Хочу разобрать подробнее» — он попадёт в страницу ответов.
Разбор внутри: хранение событий, индексы и поиск по payload.
Разбор внутри: атомарность БД-записи и отправки события.
Разбор внутри: сообщения, которые не удалось обработать после ретраев.
Разбор внутри: дубликаты, идемпотентная обработка и транзакции.
Разбор внутри: новый groupId, offsets и replay для аналитики.
Готов идти дальше?
Выбери вопросы, которые тебе интересны, и изучи их. Потом — к следующему уроку.
🚀 Перейти к Уроку 17Сначала пройди все секции и выбери хотя бы 1 вопрос
⬅️ Назад к Уроку 15