Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce
System design order dispatch untuk quick-commerce (q-commerce). Pelajari batching window untuk kumpulkan order, Hungarian/greedy matching driver ke order, Redis GEOADD/GEOSEARCH untuk cari driver terdekat, real-time driver state tracking (idle/picking/delivering), batch delivery multi-order, dan edge cases seperti driver reject, order timeout, serta starvation handling. Dilengkapi implementasi Golang dengan context.Context, error wrapping, struct-based services, channel-based event sourcing, dan concurrent pipeline.
- Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce
- Masalah: Kenapa Dispatch Q-Commerce Sulit?
- Arsitektur Dispatch Engine
- System Components
- 1. Batch Collector: Kumpulin Order, Baru Proses
- 2. Redis Geospatial: Cari Driver Terdekat
- 3. Dispatch Matcher: Greedy Assignment dengan Scoring
- 4. Driver State Machine
- 5. Re-assignment Handler: Driver Reject atau Timeout
- 6. Batch Delivery Multi-Order
- Edge Cases dan Penanganannya
- 1. Driver Reject Loop
- 2. Order Timeout
- 3. Starvation Driver Jauh
- 4. Concurrent Assignment Race
- 5. Mass Driver Drop-off
- Key Takeaways
Order Dispatch & Driver Allocation: Hungarian Algorithm untuk Q-Commerce
"Di quick-commerce, setiap detik berarti. Order yang tidak ter-assign dalam 30 detik berpotensi loss — pelanggan pergi, driver frustrasi, dan SLA hancur."
TL;DR
Sistem dispatch adalah jantung Q-Commerce. Ia harus memutuskan siapa yang mengambil order mana dalam waktu nyaris real-time dengan mempertimbangkan jarak, beban driver, ETA, dan prioritas order. Artikel ini membahas implementasi lengkap dispatch engine di Go: batching window, greedy assignment scoring, driver state machine, Redis geospatial indexing, failover mechanism, dan event sourcing untuk audit trail.
Masalah: Kenapa Dispatch Q-Commerce Sulit?
Bayangkan ini: jam makan siang, 50 order masuk dalam 2 menit. 20 driver online tersebar di berbagai lokasi. Beberapa driver sedang dalam perjalanan mengantar, beberapa baru selesai, yang lain masih menunggu di hub.
Apa yang membuat dispatch Q-Commerce berbeda dari dispatch biasa?
- Time-pressure ekstrem: Janji <30 menit dari order ke delivery. Setiap detik di dispatch pipeline adalah biaya.
- Multi-order batches: Satu driver bisa bawa 3-5 order sekali jalan — rutenya harus optimal.
- Dynamic driver state: Driver bisa reject, cancel, atau tiba-tiba offline. Sistem harus reassign dalam hitungan detik.
- Geospatial constraint: Driver hanya efisien dalam radius tertentu dari hub dan rute pengiriman.
- Starvation: Driver yang jauh jangan terus-terusan kalah oleh driver yang dekat — mereka perlu "dihangatkan" dengan prioritas.
Core Trade-off
Optimal vs Cepat. Algoritma assignment optimal (Hungarian) punya kompleksitas O(n^3). Untuk 50 order x 20 driver, perlu ~50.000 iterasi — cepat. Tapi untuk 500 order x 200 driver, Hungarian O(100.000^3) tidak feasible. Trade-off antara akurasi matching dan latency dispatch adalah keputusan arsitektural paling kritis di dispatch engine.
Arsitektur Dispatch Engine
System Components
1. Batch Collector: Kumpulin Order, Baru Proses
Tidak semua order bisa langsung di-dispatch satu per satu. Bayangkan order pertama di-assign ke driver A, lalu 3 detik kemudian muncul order yang searah dengan driver A — sayang, sudah terlanjur. Batching mengumpulkan order dalam window waktu tertentu, lalu memproses semuanya sekaligus untuk assignment yang lebih optimal.
package dispatch
import (
"context"
"fmt"
"log/slog"
"sync"
"time"
)
// BatchConfig controls batch collection behavior.
type BatchConfig struct {
MaxWindow time.Duration // max wait before flushing (e.g., 5 seconds)
MinOrders int // flush immediately if this many orders collected
MaxOrders int // hard cap per batch
}
// PendingOrder represents an order waiting to be dispatched.
type PendingOrder struct {
OrderID string
HubID string
Lat, Lng float64
CreatedAt time.Time
Priority int // higher = more urgent
}
// BatchCollector accumulates orders and triggers dispatch.
type BatchCollector struct {
cfg BatchConfig
mu sync.Mutex
batch []PendingOrder
flush chan []PendingOrder
done chan struct{}
}
// NewBatchCollector creates and starts the collector loop.
func NewBatchCollector(cfg BatchConfig) *BatchCollector {
bc := &BatchCollector{
cfg: cfg,
batch: make([]PendingOrder, 0, cfg.MaxOrders),
flush: make(chan []PendingOrder, 1),
done: make(chan struct{}),
}
go bc.loop()
return bc
}
// Enqueue adds an order to the current batch.
func (bc *BatchCollector) Enqueue(ctx context.Context, order PendingOrder) error {
bc.mu.Lock()
defer bc.mu.Unlock()
if len(bc.batch) >= bc.cfg.MaxOrders {
return ErrBatchFull
}
bc.batch = append(bc.batch, order)
if len(bc.batch) >= bc.cfg.MinOrders {
select {
case bc.flush <- bc.takeBatch():
default:
// flush already queued
}
}
return nil
}
// FlushChan returns a channel that receives batches ready for dispatch.
func (bc *BatchCollector) FlushChan() <-chan []PendingOrder {
return bc.flush
}
// Stop gracefully shuts down the collector.
func (bc *BatchCollector) Stop() {
close(bc.done)
}
func (bc *BatchCollector) loop() {
ticker := time.NewTicker(bc.cfg.MaxWindow)
defer ticker.Stop()
for {
select {
case <-ticker.C:
bc.mu.Lock()
if len(bc.batch) > 0 {
bc.flush <- bc.takeBatch()
}
bc.mu.Unlock()
case <-bc.done:
return
}
}
}
func (bc *BatchCollector) takeBatch() []PendingOrder {
batch := bc.batch
bc.batch = make([]PendingOrder, 0, bc.cfg.MaxOrders)
return batch
}
var ErrBatchFull = fmt.Errorf("batch collector: queue full")Mengapa Batching Penting?
Tanpa batching, sistem dispatch riskan terhadap sub-optimal assignment. Contoh: Order A masuk, langsung di-assign ke Driver X yang sedang dekat. 2 detik kemudian Order B masuk — ternyata lokasi Order B sangat dekat dengan Order A, dan Driver Y yang baru selesai delivery bisa ambil keduanya. Dengan batching, kita lihat A dan B bersamaan, assign ke Y, hemat 1 trip.
2. Redis Geospatial: Cari Driver Terdekat
Redis menyediakan perintah GEOADD, GEOSEARCH, dan GEODIST yang sempurna untuk query geospatial real-time. Setiap driver mengirim lokasi setiap beberapa detik (heartbeat), kita update di Redis, lalu dispatch tinggal query.
package dispatch
import (
"context"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
const (
driverLocationKey = "driver:locations"
driverStateKey = "driver:%s:state"
driverLoadKey = "driver:%s:load"
)
// DriverLocation represents a driver's current position and capacity.
type DriverLocation struct {
DriverID string `json:"driver_id"`
Lat float64 `json:"lat"`
Lng float64 `json:"lng"`
State string `json:"state"` // idle, to_hub, picking, delivering
OrderCount int `json:"order_count"` // current assigned orders
MaxOrders int `json:"max_orders"` // capacity
}
// GeoRepository wraps Redis geo operations.
type GeoRepository struct {
rdb *redis.Client
}
// NewGeoRepository creates a GeoRepository.
func NewGeoRepository(rdb *redis.Client) *GeoRepository {
return &GeoRepository{rdb: rdb}
}
// UpdateDriverLocation stores driver geo + state atomically.
func (r *GeoRepository) UpdateDriverLocation(ctx context.Context, d DriverLocation) error {
pipe := r.rdb.Pipeline()
pipe.GeoAdd(ctx, driverLocationKey, &redis.GeoLocation{
Name: d.DriverID,
Latitude: d.Lat,
Longitude: d.Lng,
})
pipe.Set(ctx, fmt.Sprintf(driverStateKey, d.DriverID), d.State, 30*time.Second)
pipe.Set(ctx, fmt.Sprintf(driverLoadKey, d.DriverID), d.OrderCount, 30*time.Second)
_, err := pipe.Exec(ctx)
return fmt.Errorf("update driver location: %w", err)
}
// FindIdleDriversNearby returns drivers in IDLE state within radius (meters).
func (r *GeoRepository) FindIdleDriversNearby(ctx context.Context, lat, lng float64, radiusMeters float64) ([]DriverLocation, error) {
res, err := r.rdb.GeoSearch(ctx, driverLocationKey, &redis.GeoSearchQuery{
Longitude: lng,
Latitude: lat,
Radius: radiusMeters,
RadiusUnit: "m",
Sort: "ASC",
Count: 50,
CountAny: false,
}).Result()
if err != nil {
return nil, fmt.Errorf("geo search: %w", err)
}
var drivers []DriverLocation
for _, name := range res {
state, err := r.rdb.Get(ctx, fmt.Sprintf(driverStateKey, name)).Result()
if err != nil || state != "idle" {
continue
}
loadStr, _ := r.rdb.Get(ctx, fmt.Sprintf(driverLoadKey, name)).Result()
load := 0
fmt.Sscanf(loadStr, "%d", &load)
drivers = append(drivers, DriverLocation{
DriverID: name,
OrderCount: load,
State: state,
})
}
return drivers, nil
}
// RemoveDriver removes a driver from geo index (e.g., when going offline).
func (r *GeoRepository) RemoveDriver(ctx context.Context, driverID string) error {
pipe := r.rdb.Pipeline()
pipe.ZRem(ctx, driverLocationKey, driverID)
pipe.Del(ctx, fmt.Sprintf(driverStateKey, driverID))
pipe.Del(ctx, fmt.Sprintf(driverLoadKey, driverID))
_, err := pipe.Exec(ctx)
return fmt.Errorf("remove driver: %w", err)
}Redis GEO Performance
GEOSEARCH menggunakan struktur data sorted set di Redis. Kompleksitas O(log N) untuk setiap anggota ditambah O(N) untuk hasil. Dengan 10.000 driver aktif, query radius 2km biasanya selesai dalam <1ms. Tapi hati-hati: Redis adalah single-threaded — GEOSEARCH dengan radius besar (seluruh kota) bisa block event loop selama puluhan ms.
3. Dispatch Matcher: Greedy Assignment dengan Scoring
Ini adalah inti sistem. Setelah batch terkumpul dan driver idle teridentifikasi, matcher menghitung score untuk setiap pasangan driver-order, lalu membuat assignment optimal.
Formula scoring:
score(d, o) = w1 * normalized_distance +
w2 * normalized_eta +
w3 * driver_load_penalty +
w4 * priority_bonus(o) +
w5 * anti_starvation(d)package dispatch
import (
"context"
"fmt"
"math"
"sort"
"sync"
)
// MatchingConfig controls scoring weights.
type MatchingConfig struct {
MaxSearchRadiusMeters float64
MaxBatchSize int
WeightDistance float64 // w1
WeightETA float64 // w2
WeightLoad float64 // w3
WeightPriority float64 // w4
WeightStarvation float64 // w5
// Anti-starvation: increment every time a driver is skipped
StarvationIncrement float64
}
// DriverMatchScore holds scoring for a driver-order pair.
type DriverMatchScore struct {
DriverID string `json:"driver_id"`
OrderID string `json:"order_id"`
Distance float64 `json:"distance_m"`
ETA float64 `json:"eta_s"`
Score float64 `json:"score"` // lower = better match
}
// MatchResult is the output of the matcher.
type MatchResult struct {
Assignments []Assignment `json:"assignments"`
Scores []DriverMatchScore `json:"scores"`
Unassigned []PendingOrder `json:"unassigned"`
}
// Assignment pairs a driver with one or more orders.
type Assignment struct {
DriverID string `json:"driver_id"`
OrderIDs []string `json:"order_ids"`
}
// Matcher performs driver-order matching.
type Matcher struct {
cfg MatchingConfig
mu sync.Mutex
starvation map[string]float64 // driver_id -> accumulated starvation score
}
// NewMatcher creates a Matcher.
func NewMatcher(cfg MatchingConfig) *Matcher {
return &Matcher{
cfg: cfg,
starvation: make(map[string]float64),
}
}
// Match assigns orders to drivers using greedy algorithm.
func (m *Matcher) Match(ctx context.Context, orders []PendingOrder, drivers []DriverLocation, geo *GeoRepository) (*MatchResult, error) {
if len(orders) == 0 || len(drivers) == 0 {
return &MatchResult{
Unassigned: orders,
}, nil
}
// Build score matrix: driver -> list of (order, score) sorted ascending
type orderScore struct {
order PendingOrder
score float64
dist float64
eta float64
}
type driverScores struct {
driver DriverLocation
scores []orderScore
}
var scoredDrivers []driverScores
for _, d := range drivers {
// Find distance and ETA via Redis or routing service
var scores []orderScore
for _, o := range orders {
dist, eta, err := m.estimateDistanceETA(ctx, geo, d, o)
if err != nil {
continue
}
score := m.computeScore(d, o, dist, eta)
scores = append(scores, orderScore{
order: o,
score: score,
dist: dist,
eta: eta,
})
}
sort.Slice(scores, func(i, j int) bool {
return scores[i].score < scores[j].score
})
scoredDrivers = append(scoredDrivers, driverScores{
driver: d,
scores: scores,
})
}
// Greedy assignment: pick lowest score pair, assign, remove from pool
assignedOrders := make(map[string]bool) // order_id
assignedDrivers := make(map[string]bool) // driver_id
var assignments []Assignment
var matchScores []DriverMatchScore
for {
bestScore := math.MaxFloat64
var bestDriverIdx int
var bestOrder orderScore
found := false
for i, ds := range scoredDrivers {
if assignedDrivers[ds.driver.DriverID] {
continue
}
for _, s := range ds.scores {
if assignedOrders[s.order.OrderID] {
continue
}
if s.score < bestScore {
bestScore = s.score
bestDriverIdx = i
bestOrder = s
found = true
}
break // first unassigned is lowest score for this driver
}
}
if !found {
break
}
driver := scoredDrivers[bestDriverIdx].driver
assignedOrders[bestOrder.order.OrderID] = true
assignedDrivers[driver.DriverID] = true
assignments = append(assignments, Assignment{
DriverID: driver.DriverID,
OrderIDs: []string{bestOrder.order.OrderID},
})
matchScores = append(matchScores, DriverMatchScore{
DriverID: driver.DriverID,
OrderID: bestOrder.order.OrderID,
Distance: bestOrder.dist,
ETA: bestOrder.eta,
Score: bestScore,
})
// Reset starvation for assigned driver
m.resetStarvation(driver.DriverID)
// Check if driver has capacity for more orders (batch delivery)
if len(assignments) > 0 && driver.OrderCount < driver.MaxOrders {
// Allow driver to take another order
delete(assignedDrivers, driver.DriverID)
}
}
// Collect unassigned orders
var unassigned []PendingOrder
for _, o := range orders {
if !assignedOrders[o.OrderID] {
m.incrementStarvationForUnassigned(scoredDrivers)
unassigned = append(unassigned, o)
}
}
return &MatchResult{
Assignments: assignments,
Scores: matchScores,
Unassigned: unassigned,
}, nil
}
func (m *Matcher) computeScore(d DriverLocation, o PendingOrder, dist, eta float64) float64 {
// Normalize distance: 0..1 based on max search radius
normDist := dist / m.cfg.MaxSearchRadiusMeters
// Normalize ETA: assume max 1800s (30 min)
normETA := eta / 1800.0
// Load penalty: driver capacity usage
loadPenalty := float64(d.OrderCount) / float64(d.MaxOrders)
starveBonus := m.getStarvation(d.DriverID)
score := m.cfg.WeightDistance*normDist +
m.cfg.WeightETA*normETA +
m.cfg.WeightLoad*loadPenalty -
m.cfg.WeightPriority*float64(o.Priority) -
m.cfg.WeightStarvation*starveBonus
return score
}
func (m *Matcher) estimateDistanceETA(ctx context.Context, geo *GeoRepository, d DriverLocation, o PendingOrder) (dist, eta float64, err error) {
// Simplified: use Haversine for distance, estimate speed
// In production: call OSRM or Google Maps Routing API
distance := haversine(d.Lat, d.Lng, o.Lat, o.Lng)
speed := 8.33 // ~30 km/h average for scooters
eta = distance / speed
return distance, eta, nil
}
func (m *Matcher) getStarvation(driverID string) float64 {
m.mu.Lock()
defer m.mu.Unlock()
return m.starvation[driverID]
}
func (m *Matcher) resetStarvation(driverID string) {
m.mu.Lock()
defer m.mu.Unlock()
delete(m.starvation, driverID)
}
func (m *Matcher) incrementStarvationForUnassigned(scoredDrivers []driverScores) {
m.mu.Lock()
defer m.mu.Unlock()
for _, ds := range scoredDrivers {
if !m.isDriverAssigned(ds.driver.DriverID) {
m.starvation[ds.driver.DriverID] += m.cfg.StarvationIncrement
}
}
}
func (m *Matcher) isDriverAssigned(driverID string) bool {
// In real impl, check against assignments
return false
}
// haversine calculates distance between two coordinates in meters.
func haversine(lat1, lng1, lat2, lng2 float64) float64 {
const R = 6371000 // Earth radius in meters
dLat := (lat2 - lat1) * math.Pi / 180
dLng := (lng2 - lng1) * math.Pi / 180
a := math.Sin(dLat/2)*math.Sin(dLat/2) +
math.Cos(lat1*math.Pi/180)*math.Cos(lat2*math.Pi/180)*
math.Sin(dLng/2)*math.Sin(dLng/2)
c := 2 * math.Atan2(math.Sqrt(a), math.Sqrt(1-a))
return R * c
}4. Driver State Machine
Setiap driver melewati siklus state yang ketat. State transition ini penting untuk:
- Visibility: Tim ops tahu persis status setiap driver
- Correctness: Sistem tahu driver mana yang available untuk dispatch
- SLA tracking: Waktu di setiap state dihitung untuk analytics
package dispatch
import (
"context"
"fmt"
"time"
)
// DriverState represents the lifecycle of a driver.
type DriverState string
const (
DriverOffline DriverState = "OFFLINE"
DriverIdle DriverState = "IDLE"
DriverAssigned DriverState = "ASSIGNED"
DriverToHub DriverState = "TO_HUB"
DriverPicking DriverState = "PICKING"
DriverDelivering DriverState = "DELIVERING"
DriverCompleted DriverState = "COMPLETED"
)
// ValidTransitions defines allowed state transitions.
var validTransitions = map[DriverState][]DriverState{
DriverOffline: {DriverIdle},
DriverIdle: {DriverAssigned, DriverOffline},
DriverAssigned: {DriverToHub, DriverIdle}, // accept or reject
DriverToHub: {DriverPicking, DriverOffline},
DriverPicking: {DriverDelivering, DriverOffline},
DriverDelivering: {DriverDelivering, DriverCompleted, DriverOffline},
DriverCompleted: {DriverIdle},
}
// StateEvent represents a driver state change event (event sourcing).
type StateEvent struct {
DriverID string `json:"driver_id"`
From DriverState `json:"from"`
To DriverState `json:"to"`
OrderID string `json:"order_id,omitempty"`
Timestamp time.Time `json:"timestamp"`
Reason string `json:"reason,omitempty"`
Coordinates []float64 `json:"coordinates,omitempty"`
}
// StateMachine manages driver lifecycle.
type StateMachine struct {
events chan StateEvent
eventSink EventStore
}
// EventStore persists state events for audit trail.
type EventStore interface {
Append(ctx context.Context, event StateEvent) error
GetByDriver(ctx context.Context, driverID string, limit int) ([]StateEvent, error)
}
// NewStateMachine creates a StateMachine.
func NewStateMachine(sink EventStore, bufferSize int) *StateMachine {
sm := &StateMachine{
events: make(chan StateEvent, bufferSize),
eventSink: sink,
}
go sm.processEvents()
return sm
}
// Transition attempts a state change and emits an event.
func (sm *StateMachine) Transition(ctx context.Context, driverID string, from, to DriverState, opts ...StateOption) error {
// Validate transition
allowed, ok := validTransitions[from]
if !ok {
return fmt.Errorf("state machine: unknown from state %q", from)
}
valid := false
for _, s := range allowed {
if s == to {
valid = true
break
}
}
if !valid {
return fmt.Errorf("state machine: invalid transition %s -> %s", from, to)
}
event := StateEvent{
DriverID: driverID,
From: from,
To: to,
Timestamp: time.Now().UTC(),
}
for _, opt := range opts {
opt(&event)
}
select {
case sm.events <- event:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
func (sm *StateMachine) processEvents() {
for event := range sm.events {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
if err := sm.eventSink.Append(ctx, event); err != nil {
// Log error, retry via DLQ
// logger.Error("failed to persist state event", "error", err)
}
cancel()
}
}
// StateOption modifies a StateEvent.
type StateOption func(*StateEvent)
// WithOrderID sets the order ID on a state event.
func WithOrderID(orderID string) StateOption {
return func(e *StateEvent) {
e.OrderID = orderID
}
}
// WithReason sets the reason for a state transition.
func WithReason(reason string) StateOption {
return func(e *StateEvent) {
e.Reason = reason
}
}
// WithCoordinates sets driver coordinates at transition time.
func WithCoordinates(lat, lng float64) StateOption {
return func(e *StateEvent) {
e.Coordinates = []float64{lat, lng}
}
}5. Re-assignment Handler: Driver Reject atau Timeout
Driver tidak selalu menerima assignment. Mereka bisa sibuk, jarak terlalu jauh, atau insentif kurang menarik. Ketika reject atau timeout terjadi, sistem harus reassign ke driver terbaik berikutnya tanpa menyebabkan starvation.
package dispatch
import (
"context"
"fmt"
"log/slog"
"time"
)
const (
driverAcceptTimeout = 15 * time.Second
maxReassignAttempts = 3
)
// ReassignmentHandler manages fallback when drivers reject.
type ReassignmentHandler struct {
matcher *Matcher
geo *GeoRepository
tracker *StateMachine
logger *slog.Logger
metrics MetricsRecorder
}
// MetricsRecorder for observability.
type MetricsRecorder interface {
IncrementCounter(name string, labels map[string]string)
ObserveHistogram(name string, value float64, labels map[string]string)
}
// NewReassignmentHandler creates a handler.
func NewReassignmentHandler(
matcher *Matcher,
geo *GeoRepository,
tracker *StateMachine,
logger *slog.Logger,
metrics MetricsRecorder,
) *ReassignmentHandler {
return &ReassignmentHandler{
matcher: matcher,
geo: geo,
tracker: tracker,
logger: logger,
metrics: metrics,
}
}
// HandleReject processes a driver rejection and finds replacement.
func (h *ReassignmentHandler) HandleReject(ctx context.Context, driverID, orderID string) error {
h.metrics.IncrementCounter("dispatch.reject", map[string]string{"driver": driverID})
// Record rejection in driver profile (for quality scoring)
if err := h.recordRejection(ctx, driverID, orderID); err != nil {
h.logger.Error("record rejection", "error", err)
}
// Return driver to IDLE
if err := h.tracker.Transition(ctx, driverID, DriverAssigned, DriverIdle,
WithReason("driver_rejected"),
); err != nil {
return fmt.Errorf("reassign: transition driver back to idle: %w", err)
}
return nil
}
// HandleTimeout processes a driver not responding to assignment.
func (h *ReassignmentHandler) HandleTimeout(ctx context.Context, driverID, orderID string) error {
h.metrics.IncrementCounter("dispatch.timeout", map[string]string{"driver": driverID})
// Penalize driver's priority for future assignments
if err := h.penalizeDriver(ctx, driverID); err != nil {
h.logger.Error("penalize driver", "error", err)
}
// Return driver to IDLE
if err := h.tracker.Transition(ctx, driverID, DriverAssigned, DriverIdle,
WithReason("response_timeout"),
); err != nil {
return fmt.Errorf("reassign: timeout transition: %w", err)
}
return nil
}
// FindReplacement finds the next best driver for an order.
func (h *ReassignmentHandler) FindReplacement(ctx context.Context, order PendingOrder, excludeDriverIDs []string) (*Assignment, error) {
for attempt := 0; attempt < maxReassignAttempts; attempt++ {
// Fetch fresh drivers
drivers, err := h.geo.FindIdleDriversNearby(ctx, order.Lat, order.Lng, 5000)
if err != nil {
return nil, fmt.Errorf("reassign find drivers: %w", err)
}
// Filter excluded
excluded := make(map[string]bool, len(excludeDriverIDs))
for _, id := range excludeDriverIDs {
excluded[id] = true
}
var available []DriverLocation
for _, d := range drivers {
if !excluded[d.DriverID] {
available = append(available, d)
}
}
if len(available) == 0 {
return nil, ErrNoDriversAvailable
}
// Re-match
result, err := h.matcher.Match(ctx, []PendingOrder{order}, available, h.geo)
if err != nil {
return nil, fmt.Errorf("reassign match: %w", err)
}
if len(result.Assignments) > 0 {
h.metrics.IncrementCounter("dispatch.reassign.success", nil)
return &result.Assignments[0], nil
}
// Exponential backoff before retry
time.Sleep(time.Duration(1<<attempt) * 500 * time.Millisecond)
}
h.metrics.IncrementCounter("dispatch.reassign.exhausted", nil)
return nil, ErrReassignExhausted
}
func (h *ReassignmentHandler) recordRejection(ctx context.Context, driverID, orderID string) error {
// Implementation: persist to DB for driver quality scoring
return nil
}
func (h *ReassignmentHandler) penalizeDriver(ctx context.Context, driverID string) error {
// Implementation: reduce driver priority score
return nil
}
// Domain errors.
var (
ErrNoDriversAvailable = fmt.Errorf("dispatch: no drivers available")
ErrReassignExhausted = fmt.Errorf("dispatch: reassign attempts exhausted")
)6. Batch Delivery Multi-Order
Satu driver bisa membawa beberapa order dalam satu trip. Ini meningkatkan efisiensi tapi menambah kompleksitas routing.
package dispatch
import (
"context"
"fmt"
"math"
"sort"
)
// BatchDelivery manages multi-order trips.
type BatchDelivery struct {
DriverID string `json:"driver_id"`
OrderIDs []string `json:"order_ids"`
Dropoffs []Dropoff `json:"dropoffs"`
OptimizedETA float64 `json:"optimized_eta_s"`
}
// Dropoff represents a delivery point.
type Dropoff struct {
OrderID string `json:"order_id"`
Lat float64 `json:"lat"`
Lng float64 `json:"lng"`
}
// RouteOptimizer optimizes multi-stop delivery route.
type RouteOptimizer struct {
geo *GeoRepository
}
// NewRouteOptimizer creates a RouteOptimizer.
func NewRouteOptimizer(geo *GeoRepository) *RouteOptimizer {
return &RouteOptimizer{geo: geo}
}
// OptimizeDropoffs returns the optimal order of drop-offs using nearest-neighbor.
func (ro *RouteOptimizer) OptimizeDropoffs(ctx context.Context, hubLat, hubLng float64, dropoffs []Dropoff) ([]Dropoff, error) {
if len(dropoffs) == 0 {
return nil, fmt.Errorf("route optimizer: empty dropoffs")
}
if len(dropoffs) == 1 {
return dropoffs, nil
}
// Nearest-neighbor heuristic from hub
sorted := make([]Dropoff, len(dropoffs))
copy(sorted, dropoffs)
currentLat, currentLng := hubLat, hubLng
for i := 0; i < len(sorted); i++ {
nearestIdx := i
nearestDist := math.MaxFloat64
for j := i; j < len(sorted); j++ {
dist := haversine(currentLat, currentLng, sorted[j].Lat, sorted[j].Lng)
if dist < nearestDist {
nearestDist = dist
nearestIdx = j
}
}
// Swap nearest to position i
sorted[i], sorted[nearestIdx] = sorted[nearestIdx], sorted[i]
currentLat, currentLng = sorted[i].Lat, sorted[i].Lng
}
return sorted, nil
}
// EstimateBatchETA calculates total ETA for a multi-order batch.
func (ro *RouteOptimizer) EstimateBatchETA(ctx context.Context, hubLat, hubLng float64, dropoffs []Dropoff) (float64, error) {
optimized, err := ro.OptimizeDropoffs(ctx, hubLat, hubLng, dropoffs)
if err != nil {
return 0, err
}
var totalETA float64
currentLat, currentLng := hubLat, hubLng
// Picking time at hub: assume 2 min per order
totalETA += float64(len(optimized)) * 120
for _, d := range optimized {
dist := haversine(currentLat, currentLng, d.Lat, d.Lng)
speed := 8.33 // m/s (~30 km/h)
totalETA += dist / speed
// Drop-off time: 30 seconds per stop
totalETA += 30
currentLat, currentLng = d.Lat, d.Lng
}
return totalETA, nil
}Edge Cases dan Penanganannya
1. Driver Reject Loop
Masalah: Seorang driver terus-menerus dapat assignment dan terus reject karena sistem tidak belajar dari pola reject sebelumnya.
Solusi: Rejection tracker dengan sliding window — jika driver reject >3x dalam 30 menit, turunkan priority score atau skip sementara.
func (h *ReassignmentHandler) shouldSkipDriver(driverID string) bool {
recentRejects, _ := h.rejectionStore.CountRecent(driverID, 30*time.Minute)
return recentRejects >= 3
}2. Order Timeout
Masalah: Order sudah di-assign ke driver, tapi driver tidak merespon dalam 15 detik. Order menunggu, pelanggan tidak sabar.
Solusi: Timer goroutine per assignment. Jika timeout, otomatis trigger reassignment dengan event timeout.
func (h *ReassignmentHandler) StartAssignmentTimer(ctx context.Context, order PendingOrder, driverID string) {
go func() {
timer := time.NewTimer(driverAcceptTimeout)
select {
case <-timer.C:
h.HandleTimeout(ctx, driverID, order.OrderID)
h.FindReplacement(ctx, order, []string{driverID})
case <-ctx.Done():
timer.Stop()
}
}()
}3. Starvation Driver Jauh
Masalah: Driver yang berada di pinggiran kota jarang dapat order karena selalu kalah skor dengan driver di pusat kota.
Solusi: Anti-starvation score incremental seperti di implementasi Matcher di atas. Setiap kali driver dilewati dalam suatu batch, starvation score naik, membuatnya lebih mungkin terpilih di batch berikutnya.
4. Concurrent Assignment Race
Masalah: Dua goroutine dispatch secara bersamaan meng-assign driver yang sama ke dua order berbeda.
Solusi: Redis locks atau database-level optimistic locking dengan versi state.
func (r *GeoRepository) TryAssignDriver(ctx context.Context, driverID string, version int) (bool, error) {
// CAS (Check-and-Set) via Redis Lua script
script := redis.NewScript(`
local state = redis.call("GET", KEYS[1])
if state == "idle" then
redis.call("SET", KEYS[1], "assigned")
return 1
end
return 0
`)
result, err := script.Run(ctx, r.rdb, []string{fmt.Sprintf(driverStateKey, driverID)}).Int()
if err != nil {
return false, err
}
return result == 1, nil
}5. Mass Driver Drop-off
Masalah: 30 driver tiba-tiba offline (gempa, demonstrasi, traffic massive). Order menumpuk.
Solusi: Dynamic radius expansion — jika idle driver dalam radius normal tidak cukup, expand radius 2x, 3x, hingga 10km. Dashboard alert ke ops team.
func findDriversWithBackoff(ctx context.Context, geo *GeoRepository, lat, lng float64) ([]DriverLocation, error) {
radii := []float64{2000, 3000, 5000, 10000}
for _, radius := range radii {
drivers, err := geo.FindIdleDriversNearby(ctx, lat, lng, radius)
if err != nil {
return nil, err
}
if len(drivers) > 0 {
return drivers, nil
}
}
return nil, ErrNoDriversAvailable
}Key Takeaways
Batch Before You Match
Jangan dispatch order satu per satu. Kumpulkan dalam batch window 3-10 detik untuk assignment yang lebih optimal. Trade-off: latency vs optimality.
Redis GEO is Your Friend
Gunakan GEOSEARCH untuk query driver terdekat dalam O(log N). Tapi jangan lupa — Redis single-threaded, jangan query radius seluruh kota.
State Machine Protects Integrity
Driver state harus strict. Validasi setiap transisi. Gunakan event sourcing untuk audit trail dan debugging.
Anti-Starvation is Design Necessity
Driver pinggiran akan terus kalah tanpa mekanisme starvation. Increment priority setiap batch dilewati. Jangan tunggu komplain driver.
Reassignment is Resilience
Driver reject atau timeout adalah normal flow, bukan error. Siapkan timer per assignment dan reassignment pipeline dengan eksponensial backoff.
Monitor Everything
Track metrics: dispatch latency, match score distribution, reject rate, timeout rate, average batch size. Grafana dashboard adalah ops tool utama.
Penutup
Dispatch engine adalah salah satu komponen paling kritikal di Q-Commerce. Kesalahan di sini langsung terasa: order tidak terkirim, pelanggan komplain, driver frustrasi. Dengan arsitektur yang tepat — batching, geospatial indexing, scoring-based assignment, strict state machine, dan anti-starvation — kamu bisa membangun sistem yang reliable bahkan di peak load 10x normal.
Kode di atas adalah production-oriented template. Untuk production, tambahkan: circuit breaker ke Redis, rate limiter di acceptance endpoint, distributed tracing (OpenTelemetry), dan chaos engineering testing untuk skenario mass driver drop-off.