Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce
Desain inventory per-hub untuk Astro (instant grocery). Reservation pattern dengan Redis DECRBY atomic + TTL, async sync ke Postgres, stock partitioning untuk hot SKU, reaper job untuk reservation leak, dan cycle counting untuk stok fisik vs sistem. Implementasi Golang lengkap.

- Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce
- Masalah: Kenapa Inventory Q-Commerce Lebih Sulit dari E-Commerce?
- Arsitektur Inventory: Redis (Operational) + Postgres (System of Record)
- Reservation Lifecycle State Machine
- 1. Inventory Service: Reservation Pattern dengan Redis Lua
- Kenapa Redis bukan Postgres?
- Redis Lua Script untuk Atomic Reserve/Release
- Edge Cases
- 2. Stock Sync Worker: Async Redis to Postgres dengan Outbox Pattern
- 3. Reaper Job: Cleanup Reservation Leak
- Problem
- Solution
- 4. Hot Key Mitigation: Stock Partitioning
- Problem
- Solution: Random Bucket Selection
- 5. Cycle Counting: Stok Fisik vs Sistem
- Problem
- Solution: Cycle Counting dengan Delta Detection
- Orchestrasi: Inventory Service
- Key Takeaways
Real-time Inventory Dark Store: Cegah Oversell di Q-Commerce
TL;DR
Oversell di q-commerce beda konsekuensinya dengan e-commerce. Kalau Tokopedia oversell, refund aja. Kalau Astro oversell Indomie, customer nunggu 30 menit — dan gak dapet. Ini bencana operational + reputasi. Artikel ini membahas reservation pattern dengan Redis atomic DECRBY + TTL, async sync ke Postgres, stock partitioning untuk hot SKU, reaper job untuk reservation leak, dan cycle counting. Kode Golang lengkap.
Masalah: Kenapa Inventory Q-Commerce Lebih Sulit dari E-Commerce?
Stok Fisik per Hub
Tiap dark store punya stok sendiri-sendiri. 100 unit Indomie di hub A, 0 di hub B. Sistem harus tahu.
Reservation Window
Dari add-to-cart sampai checkout ada jeda 5-15 menit. Stok harus di-reserve, bukan di-hold forever.
Real-time Required
Stok harus akurat dalam milidetik. Postgres eventual consistency gak cukup.
Hot SKU Contention
Indomie, Telur, Beras — 20% SKU menghasilkan 80% order. Redis key contention tinggi.
Production Horror Story
Sebuah q-commerce startup di Indonesia oversell 400% saat promo "Indomie Goreng 2 pack gratis". Penyebabnya: semua request baca stok dari Postgres replica yang lag 3 detik. 4.000 order diterima, stok cuma 1.000. 3.000 order dibatalkan. Customer trust hancur dalam 10 menit.
Arsitektur Inventory: Redis (Operational) + Postgres (System of Record)
Reservation Lifecycle State Machine
1. Inventory Service: Reservation Pattern dengan Redis Lua
Kenapa Redis bukan Postgres?
The Read/Write Ratio Problem
Di q-commerce, stok dicek 50-100x lebih sering dari pada di-update. Setiap halaman produk, setiap search, setiap cart view — semuanya baca stok. Redis bisa handle 100.000+ reads/s dengan latency < 1ms. Postgres? Di angka yang sama, CPU langsung 100% karena lock contention.
Redis Lua Script untuk Atomic Reserve/Release
package inventory
import (
"context"
"fmt"
"math/rand"
"sync"
"time"
"github.com/redis/go-redis/v9"
)
// Reservation represents a stock reservation with TTL.
type Reservation struct {
SKU string `json:"sku"`
HubID string `json:"hub_id"`
CartID string `json:"cart_id"`
UserID string `json:"user_id"`
Quantity int `json:"quantity"`
Status string `json:"status"` // RESERVED, CONFIRMED, RELEASED, EXPIRED
CreatedAt time.Time `json:"created_at"`
ExpiresAt time.Time `json:"expires_at"`
}
// InventoryConfig holds configuration for the inventory service.
type InventoryConfig struct {
DefaultReservationTTL time.Duration
MaxReservationQty int
StockKeyPrefix string
ReservationKeyPrefix string
LockKeyPrefix string
NumStockBuckets int // For hot key partitioning
}
// DefaultInventoryConfig returns sensible defaults.
func DefaultInventoryConfig() *InventoryConfig {
return &InventoryConfig{
DefaultReservationTTL: 15 * time.Minute,
MaxReservationQty: 50,
StockKeyPrefix: "inv:stock",
ReservationKeyPrefix: "inv:reservation",
LockKeyPrefix: "inv:lock",
NumStockBuckets: 10,
}
}
// InventoryReservationService handles stock reservations with Redis atomic operations.
type InventoryReservationService struct {
config *InventoryConfig
rdb *redis.Client
luaHash string
mu sync.RWMutex
}
// NewInventoryReservationService creates a new InventoryReservationService.
func NewInventoryReservationService(ctx context.Context, rdb *redis.Client, config *InventoryConfig) (*InventoryReservationService, error) {
s := &InventoryReservationService{
config: config,
rdb: rdb,
}
// Reserve Lua script: atomic check-and-decrement with TTL.
reserveScript := `
-- KEYS[1]: stock key (e.g., inv:stock:hub:sku)
-- KEYS[2]: reservation set key
-- ARGV[1]: quantity
-- ARGV[2]: reservation ID
-- ARGV[3]: TTL in seconds
-- ARGV[4]: current timestamp
-- ARGV[5]: user ID
-- ARGV[6]: cart ID
local stock = redis.call("GET", KEYS[1])
if not stock then
return {0, "SKU_NOT_FOUND"}
end
local available = tonumber(stock)
if available < tonumber(ARGV[1]) then
return {0, "INSUFFICIENT_STOCK", available}
end
-- Atomic decrement
local remaining = redis.call("DECRBY", KEYS[1], tonumber(ARGV[1]))
-- Store reservation with TTL
redis.call("HSET", KEYS[2],
"reservation_id", ARGV[2],
"sku", KEYS[1],
"quantity", ARGV[1],
"user_id", ARGV[5],
"cart_id", ARGV[6],
"status", "RESERVED",
"created_at", ARGV[4],
"expires_at", tonumber(ARGV[4]) + tonumber(ARGV[3])
)
redis.call("EXPIRE", KEYS[2], tonumber(ARGV[3]))
return {1, "RESERVED", remaining}
`
hash, err := rdb.ScriptLoad(ctx, reserveScript).Result()
if err != nil {
return nil, fmt.Errorf("inventory: load reserve lua: %w", err)
}
s.luaHash = hash
return s, nil
}
// ReserveStock atomically reserves stock for a SKU.
func (s *InventoryReservationService) ReserveStock(ctx context.Context, hubID, sku string, quantity int, cartID, userID string) (*Reservation, error) {
if quantity <= 0 || quantity > s.config.MaxReservationQty {
return nil, fmt.Errorf("inventory: invalid quantity %d (max %d)", quantity, s.config.MaxReservationQty)
}
stockKey := s.stockKey(hubID, sku)
reservationID := s.generateReservationID(cartID, sku)
reservationKey := fmt.Sprintf("%s:%s", s.config.ReservationKeyPrefix, reservationID)
ttl := int(s.config.DefaultReservationTTL.Seconds())
now := time.Now().Unix()
result, err := s.rdb.EvalSha(ctx, s.luaHash, []string{stockKey, reservationKey},
quantity, reservationID, ttl, now, userID, cartID,
).Result()
if err != nil {
return nil, fmt.Errorf("inventory: evalsha reserve: %w", err)
}
vals, ok := result.([]interface{})
if !ok || len(vals) < 2 {
return nil, fmt.Errorf("inventory: unexpected redis response: %v", result)
}
success, _ := vals[0].(int64)
if success == 0 {
msg, _ := vals[1].(string)
if len(vals) >= 3 {
available, _ := vals[2].(int64)
return nil, &InsufficientStockError{
SKU: sku,
Requested: quantity,
Available: int(available),
}
}
return nil, fmt.Errorf("inventory: reserve failed: %s", msg)
}
return &Reservation{
SKU: sku,
HubID: hubID,
CartID: cartID,
UserID: userID,
Quantity: quantity,
Status: "RESERVED",
CreatedAt: time.Now(),
ExpiresAt: time.Now().Add(s.config.DefaultReservationTTL),
}, nil
}
// ReleaseStock releases a reservation (cart abandoned, item removed).
func (s *InventoryReservationService) ReleaseStock(ctx context.Context, reservationID string) error {
releaseScript := `
-- KEYS[1]: reservation key
-- KEYS[2]: stock key
-- ARGV[1]: current timestamp
local res = redis.call("HGETALL", KEYS[1])
if #res == 0 then
return {0, "RESERVATION_NOT_FOUND"}
end
-- Extract fields
local quantity = 0
local status = ""
for i = 1, #res, 2 do
if res[i] == "quantity" then quantity = tonumber(res[i+1]) end
if res[i] == "status" then status = res[i+1] end
end
if status ~= "RESERVED" then
return {0, "WRONG_STATUS", status}
end
-- Update status
redis.call("HSET", KEYS[1], "status", "RELEASED", "released_at", ARGV[1])
-- Return stock
redis.call("INCRBY", KEYS[2], quantity)
return {1, "RELEASED", quantity}
`
result, err := s.rdb.Eval(ctx, releaseScript, []string{
fmt.Sprintf("%s:%s", s.config.ReservationKeyPrefix, reservationID),
s.stockKeyFromReservationID(reservationID),
}, time.Now().Unix()).Result()
if err != nil {
return fmt.Errorf("inventory: release stock: %w", err)
}
vals, ok := result.([]interface{})
if !ok || len(vals) < 2 {
return fmt.Errorf("inventory: unexpected release response: %v", result)
}
success, _ := vals[0].(int64)
if success == 0 {
msg, _ := vals[1].(string)
return fmt.Errorf("inventory: release failed: %s", msg)
}
return nil
}
// GetStock returns the current available stock for a SKU.
func (s *InventoryReservationService) GetStock(ctx context.Context, hubID, sku string) (int, error) {
key := s.stockKey(hubID, sku)
val, err := s.rdb.Get(ctx, key).Int()
if err != nil {
if err == redis.Nil {
return 0, nil
}
return 0, fmt.Errorf("inventory: get stock: %w", err)
}
return val, nil
}
func (s *InventoryReservationService) stockKey(hubID, sku string) string {
return fmt.Sprintf("%s:%s:%s", s.config.StockKeyPrefix, hubID, sku)
}
func (s *InventoryReservationService) generateReservationID(cartID, sku string) string {
return fmt.Sprintf("%s:%s:%d", cartID, sku, time.Now().UnixNano())
}
func (s *InventoryReservationService) stockKeyFromReservationID(reservationID string) string {
// In production, parse from the reservation hash.
// For now, derive from the reservation ID format.
return fmt.Sprintf("%s:%s", s.config.StockKeyPrefix, reservationID)
}
// InsufficientStockError is returned when there is not enough stock.
type InsufficientStockError struct {
SKU string
Requested int
Available int
}
func (e *InsufficientStockError) Error() string {
return fmt.Sprintf("inventory: insufficient stock for %s: requested %d, available %d",
e.SKU, e.Requested, e.Available)
}
// IsInsufficientStock reports whether err is an InsufficientStockError.
func IsInsufficientStock(err error) bool {
_, ok := err.(*InsufficientStockError)
return ok
}Edge Cases
Edge Case: Double Reserve
User add-to-cart item yang sama dua kali. Tanpa idempotency, stok ter-reserve 2x. Solusi: cart-level dedup. Cek apakah SKU sudah di-reserve untuk cart yang sama sebelum reserve.
Edge Case: Race Condition Confirm
User checkout tepat saat TTL reservation expired. Dua hal terjadi: reaper release stok, dan confirm reservation. Solusi: confirm harus cek status reservation di Redis (atomic) sebelum proceed.
2. Stock Sync Worker: Async Redis to Postgres dengan Outbox Pattern
Redis adalah operational store yang cepat. Postgres adalah system of record yang reliable. Keduanya harus sinkron. Tapi sinkron di setiap write? Too slow. Solusinya: outbox pattern.
package inventory
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/redis/go-redis/v9"
"github.com/segmentio/kafka-go"
)
// StockEvent represents a stock change event.
type StockEvent struct {
EventType string `json:"event_type"` // RESERVED, CONFIRMED, RELEASED, EXPIRED
HubID string `json:"hub_id"`
SKU string `json:"sku"`
Quantity int `json:"quantity"`
ReservationID string `json:"reservation_id"`
UserID string `json:"user_id"`
CartID string `json:"cart_id"`
Timestamp int64 `json:"timestamp"`
}
// StockSyncWorker asynchronously synchronizes Redis stock state to Postgres.
type StockSyncWorker struct {
rdb *redis.Client
pool *pgxpool.Pool
kafkaReader *kafka.Reader
logger *slog.Logger
batchSize int
}
// NewStockSyncWorker creates a new StockSyncWorker.
func NewStockSyncWorker(
rdb *redis.Client,
pool *pgxpool.Pool,
kafkaReader *kafka.Reader,
logger *slog.Logger,
) *StockSyncWorker {
return &StockSyncWorker{
rdb: rdb,
pool: pool,
kafkaReader: kafkaReader,
logger: logger,
batchSize: 100,
}
}
// Run starts the sync worker loop. Blocks until context is cancelled.
func (w *StockSyncWorker) Run(ctx context.Context) error {
w.logger.InfoContext(ctx, "stock sync worker started")
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
if err := w.processBatch(ctx); err != nil {
w.logger.ErrorContext(ctx, "sync worker batch failed",
"error", err,
"retry_in_seconds", 5,
)
time.Sleep(5 * time.Second)
}
}
}
}
// processBatch reads a batch of events from Kafka and syncs to Postgres.
func (w *StockSyncWorker) processBatch(ctx context.Context) error {
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
messages := make([]kafka.Message, 0, w.batchSize)
for i := 0; i < w.batchSize; i++ {
msg, err := w.kafkaReader.ReadMessage(ctx)
if err != nil {
if i > 0 {
break // Process what we have.
}
return fmt.Errorf("sync worker: read message: %w", err)
}
messages = append(messages, msg)
}
if len(messages) == 0 {
return nil
}
tx, err := w.pool.Begin(ctx)
if err != nil {
return fmt.Errorf("sync worker: begin tx: %w", err)
}
defer tx.Rollback(ctx)
for _, msg := range messages {
var event StockEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
w.logger.ErrorContext(ctx, "sync worker: unmarshal event",
"error", err,
"key", string(msg.Key),
)
continue
}
if err := w.syncEvent(ctx, tx, &event); err != nil {
return fmt.Errorf("sync worker: sync event %s: %w", event.EventType, err)
}
}
if err := tx.Commit(ctx); err != nil {
return fmt.Errorf("sync worker: commit tx: %w", err)
}
w.logger.InfoContext(ctx, "stock sync batch completed",
"events_processed", len(messages),
)
return nil
}
// syncEvent applies a stock event to Postgres.
func (w *StockSyncWorker) syncEvent(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
switch event.EventType {
case "RESERVED", "CONFIRMED":
return w.applyReservation(ctx, tx, event)
case "RELEASED", "EXPIRED":
return w.applyRelease(ctx, tx, event)
default:
return fmt.Errorf("unknown event type: %s", event.EventType)
}
}
func (w *StockSyncWorker) applyReservation(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
query := `
INSERT INTO stock_reservations (
reservation_id, hub_id, sku, user_id, cart_id,
quantity, status, created_at, expires_at
) VALUES ($1, $2, $3, $4, $5, $6, $7, NOW(), NOW() + INTERVAL '15 minutes')
ON CONFLICT (reservation_id) DO NOTHING
`
_, err := tx.Exec(ctx, query,
event.ReservationID,
event.HubID,
event.SKU,
event.UserID,
event.CartID,
event.Quantity,
event.EventType,
)
if err != nil {
return fmt.Errorf("apply reservation: %w", err)
}
// Update the inventory table.
updateQuery := `
INSERT INTO hub_inventory (hub_id, sku, reserved_quantity, updated_at)
VALUES ($1, $2, $3, NOW())
ON CONFLICT (hub_id, sku)
DO UPDATE SET reserved_quantity = hub_inventory.reserved_quantity + $3,
updated_at = NOW()
`
_, err = tx.Exec(ctx, updateQuery, event.HubID, event.SKU, event.Quantity)
if err != nil {
return fmt.Errorf("apply inventory update: %w", err)
}
return nil
}
func (w *StockSyncWorker) applyRelease(ctx context.Context, tx pgx.Tx, event *StockEvent) error {
query := `
UPDATE stock_reservations
SET status = $2, released_at = NOW()
WHERE reservation_id = $1 AND status IN ('RESERVED', 'CONFIRMED')
`
result, err := tx.Exec(ctx, query, event.ReservationID, event.EventType)
if err != nil {
return fmt.Errorf("apply release: %w", err)
}
if result.RowsAffected() == 0 {
w.logger.WarnContext(ctx, "reservation not found for release",
"reservation_id", event.ReservationID,
)
}
// Decrement reserved quantity.
updateQuery := `
UPDATE hub_inventory
SET reserved_quantity = GREATEST(0, reserved_quantity - $3),
updated_at = NOW()
WHERE hub_id = $1 AND sku = $2
`
_, err = tx.Exec(ctx, updateQuery, event.HubID, event.SKU, event.Quantity)
if err != nil {
return fmt.Errorf("apply inventory release: %w", err)
}
return nil
}
// PostgresSchema returns the DDL for the inventory tables.
// Run this during migration.
func PostgresSchema() string {
return `
CREATE TABLE IF NOT EXISTS hub_inventory (
hub_id VARCHAR(50) NOT NULL,
sku VARCHAR(100) NOT NULL,
available_quantity INTEGER NOT NULL DEFAULT 0,
reserved_quantity INTEGER NOT NULL DEFAULT 0,
pending_quantity INTEGER NOT NULL DEFAULT 0,
damaged_quantity INTEGER NOT NULL DEFAULT 0,
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (hub_id, sku)
);
CREATE TABLE IF NOT EXISTS stock_reservations (
reservation_id VARCHAR(200) PRIMARY KEY,
hub_id VARCHAR(50) NOT NULL,
sku VARCHAR(100) NOT NULL,
user_id VARCHAR(100) NOT NULL,
cart_id VARCHAR(100) NOT NULL,
quantity INTEGER NOT NULL,
status VARCHAR(20) NOT NULL DEFAULT 'RESERVED',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
expires_at TIMESTAMPTZ,
confirmed_at TIMESTAMPTZ,
released_at TIMESTAMPTZ,
cancelled_at TIMESTAMPTZ
);
CREATE INDEX idx_reservations_status ON stock_reservations(status);
CREATE INDEX idx_reservations_user ON stock_reservations(user_id);
CREATE INDEX idx_reservations_cart ON stock_reservations(cart_id);
CREATE INDEX idx_hub_inventory_updated ON hub_inventory(updated_at);
`
}3. Reaper Job: Cleanup Reservation Leak
Problem
Setiap flash sale atau promo besar, ribuan reservation dibuat. Tapi banyak yang abandoned — user add-to-cart tapi gak checkout. Tanpa cleanup, stok terkunci selamanya.
Solution
Reaper job jalan setiap 30 detik, scan reservation yang expired, release otomatis.
Reaper query Redis keys dengan TTL expired
Release stok dan update status dalam 1 atomic operation
Biar StockSyncWorker sync ke Postgres
Track release rate, leak percentage, average hold time
package inventory
import (
"context"
"fmt"
"log/slog"
"sync"
"time"
"github.com/redis/go-redis/v9"
"github.com/segmentio/kafka-go"
)
// ReaperConfig holds configuration for the reservation reaper.
type ReaperConfig struct {
// Interval between reaper runs.
Interval time.Duration
// BatchSize is the max number of reservations to process per run.
BatchSize int
// GracePeriod is how long after TTL expiry before reaping.
GracePeriod time.Duration
}
// DefaultReaperConfig returns sensible defaults.
func DefaultReaperConfig() *ReaperConfig {
return &ReaperConfig{
Interval: 30 * time.Second,
BatchSize: 1000,
GracePeriod: 30 * time.Second,
}
}
// ReaperMetrics tracks reaper performance.
type ReaperMetrics struct {
mu sync.Mutex
TotalExpired int64
TotalReleased int64
TotalErrors int64
LastRunDuration time.Duration
LastRunTime time.Time
}
// Reaper periodically cleans up expired reservations.
type Reaper struct {
config *ReaperConfig
rdb *redis.Client
kafkaWriter *kafka.Writer
logger *slog.Logger
metrics *ReaperMetrics
}
// NewReaper creates a new Reaper.
func NewReaper(config *ReaperConfig, rdb *redis.Client, kafkaWriter *kafka.Writer, logger *slog.Logger) *Reaper {
return &Reaper{
config: config,
rdb: rdb,
kafkaWriter: kafkaWriter,
logger: logger,
metrics: &ReaperMetrics{},
}
}
// Run starts the reaper loop. Blocks until context is cancelled.
func (r *Reaper) Run(ctx context.Context) error {
r.logger.InfoContext(ctx, "reaper started",
"interval", r.config.Interval,
"batch_size", r.config.BatchSize,
)
ticker := time.NewTicker(r.config.Interval)
defer ticker.Stop()
// Run immediately on start.
if err := r.reap(ctx); err != nil {
r.logger.ErrorContext(ctx, "initial reap failed", "error", err)
}
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
start := time.Now()
if err := r.reap(ctx); err != nil {
r.metrics.mu.Lock()
r.metrics.TotalErrors++
r.metrics.mu.Unlock()
r.logger.ErrorContext(ctx, "reap cycle failed", "error", err)
}
r.metrics.mu.Lock()
r.metrics.LastRunDuration = time.Since(start)
r.metrics.LastRunTime = time.Now()
r.metrics.mu.Unlock()
}
}
}
// reap performs one cycle of expired reservation cleanup.
func (r *Reaper) reap(ctx context.Context) error {
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
reapScript := `
-- KEYS[1]: reservation key or pattern
-- ARGV[1]: current timestamp
-- ARGV[2]: grace period (seconds)
-- Find expired reservations
local cursor = "0"
local expired = {}
local pattern = "inv:reservation:*"
repeat
local scan = redis.call("SCAN", cursor, "MATCH", pattern, "COUNT", 100)
cursor = scan[1]
local keys = scan[2]
for _, key in ipairs(keys) do
local expires_at = redis.call("HGET", key, "expires_at")
local status = redis.call("HGET", key, "status")
if expires_at and status == "RESERVED" then
local expiry = tonumber(expires_at)
if tonumber(ARGV[1]) >= (expiry + tonumber(ARGV[2])) then
-- This reservation is expired
local quantity = redis.call("HGET", key, "quantity")
local sku = redis.call("HGET", key, "sku")
-- Release stock
redis.call("INCRBY", sku, quantity)
-- Update status
redis.call("HSET", key, "status", "EXPIRED", "released_at", ARGV[1])
redis.call("EXPIRE", key, 86400) -- Keep for 24h audit
table.insert(expired, key)
table.insert(expired, quantity)
end
end
end
until cursor == "0"
return expired
`
now := time.Now().Unix()
graceSec := int(r.config.GracePeriod.Seconds())
result, err := r.rdb.Eval(ctx, reapScript, []string{}, now, graceSec).Result()
if err != nil {
return fmt.Errorf("reaper: eval reap script: %w", err)
}
expiredList, ok := result.([]interface{})
if !ok {
return nil
}
releasedCount := len(expiredList) / 2
if releasedCount == 0 {
return nil
}
// Publish release events to Kafka for Postgres sync.
for i := 0; i < len(expiredList); i += 2 {
reservationKey, _ := expiredList[i].(string)
quantity, _ := expiredList[i+1].(int64)
event := StockEvent{
EventType: "EXPIRED",
ReservationID: extractReservationID(reservationKey),
Quantity: int(quantity),
Timestamp: now,
}
eventBytes, _ := json.Marshal(event)
err := r.kafkaWriter.WriteMessages(ctx, kafka.Message{
Key: []byte(event.ReservationID),
Value: eventBytes,
})
if err != nil {
r.logger.ErrorContext(ctx, "reaper: publish event failed",
"error", err,
"reservation", event.ReservationID,
)
}
}
r.metrics.mu.Lock()
r.metrics.TotalExpired += int64(len(expiredList))
r.metrics.TotalReleased += int64(releasedCount)
r.metrics.mu.Unlock()
r.logger.InfoContext(ctx, "reaper cycle completed",
"reservations_released", releasedCount,
)
return nil
}
// Metrics returns a copy of the current reaper metrics.
func (r *Reaper) Metrics() ReaperMetrics {
r.metrics.mu.Lock()
defer r.metrics.mu.Unlock()
return ReaperMetrics{
TotalExpired: r.metrics.TotalExpired,
TotalReleased: r.metrics.TotalReleased,
TotalErrors: r.metrics.TotalErrors,
LastRunDuration: r.metrics.LastRunDuration,
LastRunTime: r.metrics.LastRunTime,
}
}
func extractReservationID(key string) string {
// Format: inv:reservation:<cart_id>:<sku>:<timestamp>
prefix := "inv:reservation:"
if len(key) > len(prefix) {
return key[len(prefix):]
}
return key
}4. Hot Key Mitigation: Stock Partitioning
Problem
SKU populer ("Indomie Goreng", "Telur 1kg", "Beras 5kg") adalah hot key. 20% SKU menghasilkan 80% traffic. Semua request tulis ke Redis key yang sama menyebabkan contention.
Solution: Random Bucket Selection
package inventory
import (
"context"
"fmt"
"math/rand"
"time"
)
// StockPartitioner distributes hot SKU stock across multiple Redis keys.
type StockPartitioner struct {
numBuckets int
prefix string
}
// NewStockPartitioner creates a new StockPartitioner.
func NewStockPartitioner(numBuckets int) *StockPartitioner {
if numBuckets <= 0 {
numBuckets = 10
}
return &StockPartitioner{
numBuckets: numBuckets,
prefix: "inv:partition",
}
}
// PartitionKey returns a specific bucket key for a hub/SKU combination.
func (p *StockPartitioner) PartitionKey(hubID, sku string, bucketIdx int) string {
return fmt.Sprintf("%s:%s:%s:%d", p.prefix, hubID, sku, bucketIdx)
}
// SelectBucket deterministically selects a bucket for a hub/SKU, with random
// fallback to spread load across buckets.
func (p *StockPartitioner) SelectBucket(hubID, sku string, useRandom bool) int {
if useRandom {
return rand.Intn(p.numBuckets)
}
// Deterministic hash for consistent bucket selection per SKU.
// This reduces cache misses since same SKU tends to hit same bucket.
hash := hashString(hubID + ":" + sku)
return hash % p.numBuckets
}
// PartitionedReserve reserves stock across partitioned buckets with fallback.
func (p *StockPartitioner) PartitionedReserve(
svc *InventoryReservationService,
ctx context.Context,
hubID, sku string,
quantity int,
cartID, userID string,
) (*Reservation, error) {
// Try primary bucket first, then fallback to random.
for attempt := 0; attempt < p.numBuckets; attempt++ {
bucketIdx := p.SelectBucket(sku, hubID, attempt > 0)
partitionedSKU := fmt.Sprintf("%s:%d", sku, bucketIdx)
reservation, err := svc.ReserveStock(ctx, hubID, partitionedSKU, quantity, cartID, userID)
if err == nil {
return reservation, nil
}
if !IsInsufficientStock(err) {
return nil, err // Real error, not stock issue.
}
// Try next bucket.
continue
}
return nil, &InsufficientStockError{
SKU: sku,
Requested: quantity,
Available: 0,
}
}
// InitializePartitionedStock initializes partitioned stock for a hub/SKU.
func (p *StockPartitioner) InitializePartitionedStock(
ctx context.Context,
rdb *redis.Client,
hubID, sku string,
totalStock int,
ttl time.Duration,
) error {
stockPerBucket := totalStock / p.numBuckets
remainder := totalStock % p.numBuckets
pipe := rdb.Pipeline()
for i := 0; i < p.numBuckets; i++ {
bucketStock := stockPerBucket
if i < remainder {
bucketStock++
}
key := p.PartitionKey(hubID, sku, i)
pipe.Set(ctx, key, bucketStock, ttl)
}
_, err := pipe.Exec(ctx)
if err != nil {
return fmt.Errorf("stock partitioner: init stock: %w", err)
}
return nil
}
func hashString(s string) int {
h := 0
for _, c := range s {
h = h*31 + int(c)
}
if h < 0 {
h = -h
}
return h
}5. Cycle Counting: Stok Fisik vs Sistem
Problem
Stok di sistem dan stok di dark store sering berbeda. Picker ambil barang dari rak → stok sistem gak update. Barang rusak → gak dicatat. Stok hilang (shrinkage) → sistem tahu-nya masih ada.
Solution: Cycle Counting dengan Delta Detection
package inventory
import (
"context"
"fmt"
"log/slog"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
// CycleCountService handles physical vs system stock reconciliation.
type CycleCountService struct {
pool *pgxpool.Pool
logger *slog.Logger
}
// NewCycleCountService creates a new CycleCountService.
func NewCycleCountService(pool *pgxpool.Pool, logger *slog.Logger) *CycleCountService {
return &CycleCountService{
pool: pool,
logger: logger,
}
}
// CountResult represents the result of a cycle count for a single SKU.
type CountResult struct {
HubID string
SKU string
PhysicalStock int
SystemStock int
ReservedStock int
AvailableStock int // system=available - reserved
Delta int
DeltaPercent float64
CountedAt time.Time
NeedsAdjustment bool
}
// RecordCycleCount records a physical cycle count and detects discrepancies.
func (s *CycleCountService) RecordCycleCount(ctx context.Context, hubID, sku string, physicalCount int, countedBy string) (*CountResult, error) {
// Get current system stock from Postgres.
var systemAvailable, systemReserved int
err := s.pool.QueryRow(ctx, `
SELECT available_quantity, reserved_quantity
FROM hub_inventory
WHERE hub_id = $1 AND sku = $2
`, hubID, sku).Scan(&systemAvailable, &systemReserved)
if err != nil {
return nil, fmt.Errorf("cycle count: get system stock: %w", err)
}
effectiveSystem := systemAvailable - systemReserved
delta := physicalCount - effectiveSystem
deltaPct := 0.0
if effectiveSystem > 0 {
deltaPct = float64(delta) / float64(effectiveSystem) * 100
}
result := &CountResult{
HubID: hubID,
SKU: sku,
PhysicalStock: physicalCount,
SystemStock: systemAvailable,
ReservedStock: systemReserved,
AvailableStock: effectiveSystem,
Delta: delta,
DeltaPercent: deltaPct,
CountedAt: time.Now(),
NeedsAdjustment: delta != 0,
}
// Record the count in the audit table.
_, err = s.pool.Exec(ctx, `
INSERT INTO cycle_counts (
hub_id, sku, physical_quantity, system_quantity,
delta, counted_by, counted_at
) VALUES ($1, $2, $3, $4, $5, $6, NOW())
`, hubID, sku, physicalCount, effectiveSystem, delta, countedBy)
if err != nil {
return nil, fmt.Errorf("cycle count: insert record: %w", err)
}
// If delta exceeds threshold, flag for adjustment.
if deltaPct > 5.0 || deltaPct < -5.0 {
_, err = s.pool.Exec(ctx, `
INSERT INTO inventory_adjustments (
hub_id, sku, previous_quantity, adjusted_quantity,
delta, reason, status, created_at
) VALUES ($1, $2, $3, $4, $5, $6, 'PENDING', NOW())
`, hubID, sku, effectiveSystem, physicalCount, delta, "cycle_count_discrepancy")
if err != nil {
return nil, fmt.Errorf("cycle count: create adjustment: %w", err)
}
s.logger.InfoContext(ctx, "cycle count discrepancy flagged",
"hub_id", hubID,
"sku", sku,
"delta", delta,
"delta_pct", deltaPct,
)
}
return result, nil
}
// GetDiscrepancyReport returns all SKUs with significant discrepancies.
func (s *CycleCountService) GetDiscrepancyReport(ctx context.Context, hubID string, thresholdPct float64) ([]*CountResult, error) {
rows, err := s.pool.Query(ctx, `
SELECT
cc.hub_id, cc.sku, cc.physical_quantity,
hi.available_quantity, hi.reserved_quantity,
cc.delta, cc.counted_at
FROM cycle_counts cc
JOIN hub_inventory hi ON hi.hub_id = cc.hub_id AND hi.sku = cc.sku
WHERE cc.hub_id = $1
AND ABS(cc.delta::float / NULLIF(hi.available_quantity - hi.reserved_quantity, 0)) * 100 > $2
AND cc.counted_at > NOW() - INTERVAL '24 hours'
ORDER BY ABS(cc.delta) DESC
`, hubID, thresholdPct)
if err != nil {
return nil, fmt.Errorf("cycle count: get discrepancy report: %w", err)
}
defer rows.Close()
var results []*CountResult
for rows.Next() {
var r CountResult
err := rows.Scan(&r.HubID, &r.SKU, &r.PhysicalStock, &r.SystemStock,
&r.ReservedStock, &r.Delta, &r.CountedAt)
if err != nil {
return nil, fmt.Errorf("cycle count: scan result: %w", err)
}
r.AvailableStock = r.SystemStock - r.ReservedStock
if r.AvailableStock > 0 {
r.DeltaPercent = float64(r.Delta) / float64(r.AvailableStock) * 100
}
r.NeedsAdjustment = true
results = append(results, &r)
}
return results, nil
}
// CycleCountSchema returns DDL for cycle counting tables.
func CycleCountSchema() string {
return `
CREATE TABLE IF NOT EXISTS cycle_counts (
id BIGSERIAL PRIMARY KEY,
hub_id VARCHAR(50) NOT NULL,
sku VARCHAR(100) NOT NULL,
physical_quantity INTEGER NOT NULL,
system_quantity INTEGER NOT NULL,
delta INTEGER NOT NULL,
counted_by VARCHAR(100) NOT NULL,
counted_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
notes TEXT
);
CREATE INDEX idx_cycle_counts_hub ON cycle_counts(hub_id, counted_at);
CREATE TABLE IF NOT EXISTS inventory_adjustments (
id BIGSERIAL PRIMARY KEY,
hub_id VARCHAR(50) NOT NULL,
sku VARCHAR(100) NOT NULL,
previous_quantity INTEGER NOT NULL,
adjusted_quantity INTEGER NOT NULL,
delta INTEGER NOT NULL,
reason VARCHAR(200) NOT NULL,
status VARCHAR(20) NOT NULL DEFAULT 'PENDING',
approved_by VARCHAR(100),
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
approved_at TIMESTAMPTZ
);
CREATE INDEX idx_adjustments_status ON inventory_adjustments(status);
`
}Orchestrasi: Inventory Service
package inventory
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"net/http"
"time"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/redis/go-redis/v9"
"github.com/segmentio/kafka-go"
)
// InventoryService orchestrates all inventory operations.
type InventoryService struct {
reservation *InventoryReservationService
partitioner *StockPartitioner
reaper *Reaper
syncWorker *StockSyncWorker
cycleCount *CycleCountService
rdb *redis.Client
pool *pgxpool.Pool
kafkaWriter *kafka.Writer
logger *slog.Logger
}
// NewInventoryService creates a fully wired inventory service.
func NewInventoryService(
ctx context.Context,
rdb *redis.Client,
pool *pgxpool.Pool,
logger *slog.Logger,
) (*InventoryService, error) {
config := DefaultInventoryConfig()
reservation, err := NewInventoryReservationService(ctx, rdb, config)
if err != nil {
return nil, fmt.Errorf("inventory: init reservation: %w", err)
}
kafkaWriter := &kafka.Writer{
Addr: kafka.TCP("localhost:9092"),
Topic: "inventory.events",
Balancer: &kafka.Hash{},
Async: true,
}
kafkaReader := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "inventory.events",
GroupID: "inventory-sync-worker",
MinBytes: 10,
MaxBytes: 10e6,
MaxWait: 1 * time.Second,
})
reaper := NewReaper(DefaultReaperConfig(), rdb, kafkaWriter, logger)
syncWorker := NewStockSyncWorker(rdb, pool, kafkaReader, logger)
cycleCount := NewCycleCountService(pool, logger)
return &InventoryService{
reservation: reservation,
partitioner: NewStockPartitioner(10),
reaper: reaper,
syncWorker: syncWorker,
cycleCount: cycleCount,
rdb: rdb,
pool: pool,
kafkaWriter: kafkaWriter,
logger: logger,
}, nil
}
// Start starts all background workers.
func (s *InventoryService) Start(ctx context.Context) {
go func() {
if err := s.syncWorker.Run(ctx); err != nil && err != context.Canceled {
s.logger.ErrorContext(ctx, "sync worker exited", "error", err)
}
}()
go func() {
if err := s.reaper.Run(ctx); err != nil && err != context.Canceled {
s.logger.ErrorContext(ctx, "reaper exited", "error", err)
}
}()
}
// ReserveForCart handles cart-level reservation with hot key mitigation.
func (s *InventoryService) ReserveForCart(ctx context.Context, hubID string, items []CartItem, cartID, userID string) ([]*Reservation, error) {
reservations := make([]*Reservation, 0, len(items))
for _, item := range items {
reservation, err := s.partitioner.PartitionedReserve(
s.reservation, ctx, hubID, item.SKU, item.Quantity, cartID, userID,
)
if err != nil {
// Rollback previous reservations on failure.
for _, r := range reservations {
if releaseErr := s.reservation.ReleaseStock(ctx, fmt.Sprintf("%s:%s:%d", cartID, r.SKU, r.CreatedAt.UnixNano())); releaseErr != nil {
s.logger.ErrorContext(ctx, "rollback reservation failed",
"error", releaseErr,
"sku", r.SKU,
)
}
}
return nil, fmt.Errorf("inventory: reserve cart item %s: %w", item.SKU, err)
}
reservations = append(reservations, reservation)
}
return reservations, nil
}
// CartItem represents an item in a shopping cart.
type CartItem struct {
SKU string `json:"sku"`
Quantity int `json:"quantity"`
}
// HTTPHandler returns HTTP handlers for the inventory service.
func (s *InventoryService) HTTPHandler() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("GET /inventory/stock/{hub_id}/{sku}", func(w http.ResponseWriter, r *http.Request) {
hubID := r.PathValue("hub_id")
sku := r.PathValue("sku")
stock, err := s.reservation.GetStock(r.Context(), hubID, sku)
if err != nil {
http.Error(w, `{"error":"internal_error"}`, http.StatusInternalServerError)
return
}
json.NewEncoder(w).Encode(map[string]interface{}{
"hub_id": hubID,
"sku": sku,
"stock": stock,
})
})
mux.HandleFunc("POST /inventory/reserve", func(w http.ResponseWriter, r *http.Request) {
var req struct {
HubID string `json:"hub_id"`
CartID string `json:"cart_id"`
UserID string `json:"user_id"`
Items []CartItem `json:"items"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, `{"error":"invalid_request"}`, http.StatusBadRequest)
return
}
reservations, err := s.ReserveForCart(r.Context(), req.HubID, req.Items, req.CartID, req.UserID)
if err != nil {
http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusConflict)
return
}
json.NewEncoder(w).Encode(map[string]interface{}{
"status": "reserved",
"reservations": reservations,
})
})
mux.HandleFunc("POST /inventory/release/{reservation_id}", func(w http.ResponseWriter, r *http.Request) {
reservationID := r.PathValue("reservation_id")
if err := s.reservation.ReleaseStock(r.Context(), reservationID); err != nil {
http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusBadRequest)
return
}
json.NewEncoder(w).Encode(map[string]string{"status": "released"})
})
mux.HandleFunc("POST /inventory/cycle-count", func(w http.ResponseWriter, r *http.Request) {
var req struct {
HubID string `json:"hub_id"`
SKU string `json:"sku"`
PhysicalQty int `json:"physical_quantity"`
CountedBy string `json:"counted_by"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
http.Error(w, `{"error":"invalid_request"}`, http.StatusBadRequest)
return
}
result, err := s.cycleCount.RecordCycleCount(r.Context(), req.HubID, req.SKU, req.PhysicalQty, req.CountedBy)
if err != nil {
http.Error(w, `{"error":"`+err.Error()+`"}`, http.StatusInternalServerError)
return
}
json.NewEncoder(w).Encode(result)
})
return mux
}Key Takeaways
Redis Atomic Operations
Gunakan Redis Lua script untuk check-and-decrement atomic. Single key contention dihindari dengan stock partitioning.
TTL-based Auto Release
Setiap reservation punya TTL 15 menit. Reaper job otomatis cleanup expired reservation tiap 30 detik.
Async Sync ke Postgres
Outbox pattern via Kafka. Redis untuk operational (ms-latency), Postgres untuk system of record (reliable).
Stock Partitioning
Hot SKU di-partition ke 10 bucket random. Kurangi contention di key yang sama.
Cycle Counting
Stok fisik vs sistem di-reconcile via cycle counting. Delta >5% auto-flag buat adjustment.
Idempotency & Rollback
Cart-level idempotency cegah double reserve. Partial failure trigger rollback semua reservation.
Bottom Line
Inventory q-commerce bukan soal "stok berapa" — tapi soal "stok ini sudah di-reserve siapa, berapa lama, dan apa yang terjadi kalau gak jadi dibeli?" Kombinasi Redis (fast + atomic) + Postgres (reliable + queryable) + Reaper (auto-cleanup) + Cycle Count (detect discrepancy) adalah fondasi yang sudah terbukti di Astro, Instacart, Gorillas, dan Getir.