Chuyển tới nội dung chính

Message queue & outbox

govex-cloud-mq là hạ tầng messaging dùng chung: publish/consume event giữa các service theo outbox pattern trên database, có kiểm tra idempotent ở phía consumer, gửi lại bản tin lỗi (relay) và API EventPublisher không phụ thuộc transport (Spring Event trong cùng JVM hoặc Kafka).

Khi nào sử dụng

  • Service cần phát event nghiệp vụ cho service khác consume event mà không muốn gắn chặt vào Kafka.
  • Cần đảm bảo bản tin không bị mất khi transaction nghiệp vụ rollback (outbox pattern).
  • Cần chống xử lý trùng bản tin ở consumer (idempotent theo messageId).
  • Cần cơ chế gửi lại bản tin chưa thành công và tra cứu dead-letter khi vận hành.

Cài đặt

Thêm dependency vào pom.xml — không khai version vì đã được quản lý qua BOM:

<dependency>
<groupId>vn.govex.cloud</groupId>
<artifactId>govex-cloud-mq</artifactId>
</dependency>

Nếu ứng dụng chưa dùng parent/BOM của Govex Cloud, xem hướng dẫn tại Cài đặt.

Cấu hình

Toàn bộ cấu hình nằm dưới prefix govex.mq:

PropertyMô tảMặc định
govex.mq.enabledBật auto-configuration của dependency.true
govex.mq.transportTransport cho EventPublisher: spring (in-process) hoặc kafka.spring
govex.mq.app-codeMã ứng dụng; bắt buộc khi relay.enabled=true.
govex.mq.message-tableBảng outbox lưu bản tin chờ gửi.message
govex.mq.message-tracking-tableBảng tracking trạng thái xử lý theo consumer.message_tracking
govex.mq.schema-update-enabledTự cập nhật schema hai bảng trên khi khởi động.false
govex.mq.idempotent.enabledBật kiểm tra idempotent cho consumer.true
govex.mq.idempotent.max-retriesSố lần xử lý lại tối đa cho một bản tin.3
govex.mq.relay.enabledBật relay gửi lại bản tin chưa thành công.false
govex.mq.relay.batch-sizeSố bản tin tối đa mỗi lô relay.100
govex.mq.relay.max-retriesSố lần gửi lại tối đa trước khi vào dead-letter.3
govex.mq.relay.max-retry-delayTrần backoff (ms) giữa các lần gửi lại.1000
govex.mq.relay.schedule-trigger.trigger-typeKiểu lịch relay: CRON, FIXED_RATE, FIXED_DELAY, RUN_ONCE.FIXED_RATE
govex.mq.relay.schedule-trigger.fixed-rateChu kỳ chạy relay (Duration, ví dụ 5s); dùng khi trigger-type=FIXED_RATE.
govex.mq.relay.schedule-trigger.cron-expressionBiểu thức cron; dùng khi trigger-type=CRON.0 */10 * * * *
govex.mq.kafka-restful-enabledBật endpoint POST /kafka/dispatch để đẩy bản tin vào listener qua HTTP.false
govex.mq.kafka-template-proxy-enabledBật proxy KafkaTemplate để các lệnh send đi qua outbox.false

Sử dụng

Cấu hình Kafka kèm relay:

govex:
mq:
enabled: true
transport: kafka
app-code: order-service
relay:
enabled: true
batch-size: 100
max-retries: 3
max-retry-delay: 1000
schedule-trigger:
trigger-type: FIXED_RATE
fixed-rate: 5s
idempotent:
enabled: true
max-retries: 3

Publish event không phụ thuộc transport — chỉ cần inject EventPublisher:

@Service
@RequiredArgsConstructor
public class OrderService {

private final EventPublisher eventPublisher;

@Transactional
public void createOrder(Order order) {
// ghi outbox trong cùng transaction; transport gửi sau khi commit
eventPublisher.publish("order.created", order.getId(), order,
Map.of("source", "order-service"));
}
}
  • publish(...) gửi theo transport đang cấu hình; asyncPublish(topic, payload) chỉ ghi outbox để relay gửi sau (với Kafka) hoặc fire-and-forget (với Spring Event).
  • Với transport=spring không cần Kafka: sự kiện được phát nội bộ trong cùng JVM.

Consumer dùng annotation chung cho cả hai transport:

@Component
public class OrderListener {

@MessageListener(topic = "order.created", group = "order-service")
public void onOrderCreated(OrderEvent payload) {
// xử lý event
}
}

Consumer Kafka có kiểm tra idempotent — messageId bắt buộc, hỗ trợ SpEL:

@KafkaMessageListener(
topics = "order.created",
groupId = "order-service",
messageId = "#{payload.messageId}",
traceId = "#{payload.traceId}")
public void onOrderCreated(OrderEvent payload) {
// chỉ xử lý một lần cho mỗi messageId
}

Vận hành: xem dead-letter qua GET /kafka/dead-letter (cần relay.enabled=true); đẩy bản tin vào listener qua POST /kafka/dispatch?topic=... (cần kafka-restful-enabled=true).

Luồng hoạt động

Sơ đồ dưới đây mô tả luồng outbox từ lúc publish tới consumer:

Lưu ý

  • Cần tạo bảng message, message_tracking (và received_messages nếu dùng dispatcher đồng bộ) trước khi chạy, hoặc bật schema-update-enabled=true để dependency tự cập nhật schema hai bảng message khi khởi động.
  • Relay yêu cầu Redis/Redisson cho lock giữa các instance; tính năng gửi thật qua Kafka yêu cầu spring-kafka trên classpath và cấu hình Kafka của ứng dụng.
  • Bật relay.enabled=true mà thiếu app-code sẽ làm ứng dụng khởi động thất bại ngay.
  • messageId trong @KafkaMessageListener không có giá trị mặc định — annotation báo lỗi biên dịch/runtime nếu bỏ trống.
  • Nếu ứng dụng đã có bean ErrorDecoder/KafkaTemplate tùy biến, cần kiểm tra tương thích trước khi bật kafka-template-proxy-enabled.