УРОК 16 / 30 0%

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

  1. 1
    Flyway-миграция V4__add_task_events.sql
    
    CREATE 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. 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. 3
    repository/TaskEventLogRepository.java
    
    package 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. 4
    Расширяем event/TaskCreatedEvent — добавляем snapshot
    
    public record TaskCreatedEvent(
        Long taskId,
        String title,
        String description,
        boolean done,
        String ownerUsername,
        Instant occurredAt
    ) implements TaskEvent {
        @Override public String eventType() { return "TaskCreated"; }
    }
    
  5. 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. 6
    service/AuditService.java — отдаём audit log
    
    package 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 List getTaskHistory(Long taskId) {
            return repo.findByTaskIdOrderByOccurredAtDesc(taskId)
                .stream()
                .map(AuditEventResponse::from)
                .toList();
        }
    }
    
  7. 7
    DTO
    
    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());
        }
    }
    
  8. 8
    Controller
    
    @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 List getHistory(@PathVariable Long taskId) {
            return auditService.getTaskHistory(taskId);
        }
    }
    
  9. 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;"
    

    ---

🎯

Зачем это на собесе

После урока ты должен уметь ответить на:

GDPR / SOX / HIPAA — compliance требует immutable audit log. Kafka + append-only table даёт это бесплатно.
JSONB индексируется в PostgreSQL — можно искать по содержимому события (быстро).
DLQ (Dead Letter Queue) — если событие не парсится, уходит в task-events-dlq для ручной обработки.
Outbox pattern — для exactly-once delivery (упомяни на собесе).

5 вопросов на углубление

Раскрой вопрос и нажми «🤔 Хочу разобрать подробнее» — он попадёт в страницу ответов.

1
JSONB vs JSON/text в PostgreSQL для event log

Разбор внутри: хранение событий, индексы и поиск по payload.

2
Что такое Outbox pattern и какую проблему он решает?

Разбор внутри: атомарность БД-записи и отправки события.

3
Зачем нужна DLQ в event-driven системе?

Разбор внутри: сообщения, которые не удалось обработать после ретраев.

4
At-least-once vs exactly-once — что реально гарантирует Kafka?

Разбор внутри: дубликаты, идемпотентная обработка и транзакции.

5
Как перечитать события заново без поломки текущего consumer-а?

Разбор внутри: новый groupId, offsets и replay для аналитики.

Готов идти дальше?

Выбери вопросы, которые тебе интересны, и изучи их. Потом — к следующему уроку.

🚀 Перейти к Уроку 17

Сначала пройди все секции и выбери хотя бы 1 вопрос

⬅️ Назад к Уроку 15
🎉
Новый тир
Новый тир достигнут!