Business Rules Engine - HLD
Difficulty: IntermediateβAdvanced Prerequisites: Message Queues, Caching Asked at: Amazon, Flipkart, Razorpay, PayPal, Uber, Goldman Sachs
TL;DR
A business rules engine decouples decision logic from application code. Product teams author rules (βif order > $500 AND user is new, apply 10% discountβ) via a management UI. The engine evaluates incoming events against thousands of active rules in sub-50ms, then triggers actions without redeploying any service.
flowchart LR
ADMIN["Rule Authors<br/>create and version rules"]:::client
API["Rules API"]:::service
DB[("Rule Store<br/>Postgres")]:::data
CACHE[("Rule Cache<br/>Redis + In-Memory")]:::data
ENGINE["Evaluation Engine<br/>Rete network"]:::service
KAFKA["Action Bus<br/>Kafka"]:::async
ACTIONS["Action Executor<br/>webhook email DB"]:::service
ADMIN -->|"1. Create rule with AST"| API
API -->|"2. Store versioned rule"| DB
DB -->|"3. Propagate change event"| KAFKA
KAFKA -->|"4. Invalidate and reload"| CACHE
ENGINE -->|"5. Load compiled rules"| CACHE
ENGINE -->|"6. Publish matched actions"| KAFKA
KAFKA -->|"7. Execute actions"| ACTIONS
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
| Color | Role |
|---|---|
| Purple (dark) | Client / Admin |
| Green | Service |
| Purple (bright) | Async / broker |
| Yellow | Data store |
In 3 sentences: Rule authors define conditions as structured expressions (stored as pre-parsed ASTs in JSONB). The evaluation engine loads compiled rules from in-memory cache, runs them through a Rete network for shared sub-expression optimization, and collects matches in < 50ms. Matched actions are dispatched via Kafka to a decoupled executor with retries and DLQs.
Functional Requirements
Core (top 3)
- Author and manage rules β CRUD rules with structured conditions (comparisons, arithmetic, boolean logic, nested expressions), assign actions, activate/deactivate, version with rollback.
- Evaluate events against rules β given an event payload (JSON), find all matching rules and return decisions + trigger actions, in < 50ms P99 for 10K active rules.
- Hot reload without downtime β when a rule changes, all engine instances pick up the new version within 2 seconds without restart or dropped evaluations.
Below the line
- Visual rule builder UI (assume it calls our API).
- A/B testing of rule variants.
- Rule simulation across historical event datasets.
- Multi-tenant isolation at the compute level (we do namespace isolation only).
- Complex event processing (time-windowed correlations across multiple events).
Non-Functional Requirements
- Latency β P99 evaluation < 50ms for 10K active rules against a single event.
- Throughput β 50K evaluations/sec per engine instance, horizontally scalable.
- Consistency β rule changes propagated to all instances within 2s. Strong consistency for blocking rules via direct DB read.
- Auditability β full immutable execution log for compliance.
Core Entities
- Rule β named decision unit: identity, namespace, priority, status (DRAFT/ACTIVE/DISABLED/ARCHIVED), pointer to current version.
- Rule Version β immutable snapshot of condition AST + actions. Full history for rollback.
- Condition (AST Node) β tree node representing one logic piece. Recursively composed.
- Action β what fires on match: webhook, email, field enrichment, block/allow, DB update.
- Event β incoming payload to evaluate: JSON with type and arbitrary fields.
- Namespace β logical rule grouping (e.g.,
fraud-detection,pricing). Evaluated together. - Execution Record β audit trail: event ID, matched rules, actions, latency.
API / System Interface
POST /v1/rules -> Rule
GET /v1/rules/:id -> Rule + current version
PUT /v1/rules/:id -> new version created
DELETE /v1/rules/:id -> soft delete (ARCHIVED)
POST /v1/rules/:id/activate -> set status ACTIVE
POST /v1/rules/:id/deactivate -> set status DISABLED
POST /v1/rules/:id/rollback/:version -> revert to prior version
GET /v1/rules/:id/versions -> paginated version history
POST /v1/evaluate -> EvaluationResult
POST /v1/evaluate/dry-run -> EvaluationResult (no actions fired)
GET /v1/executions?event_id=X -> execution audit records
Example β create a rule:
POST /v1/rules
{
"namespace": "fraud-detection",
"name": "high-value-new-user-block",
"condition": "order.total > 500 AND user.tier == 'NEW' AND user.country IN ['NG','GH']",
"actions": [
{ "type": "BLOCK_TRANSACTION", "reason": "High value new user from high-risk region" },
{ "type": "WEBHOOK", "url": "https://fraud-team.internal/review" }
],
"priority": 100
}
The API parses the condition string into an AST at write time (ANTLR), validates it, stores the AST in rule_versions.condition_ast.
Example β evaluate:
POST /v1/evaluate
{
"namespace": "fraud-detection",
"event_type": "order.created",
"event_id": "evt_abc123",
"payload": {
"order": { "total": 750, "currency": "USD" },
"user": { "id": "u_99", "tier": "NEW", "country": "NG" }
}
}
// Response: { "evaluation_time_ms": 12.4, "matched_rules": [...], "final_decision": "BLOCK" }
Security notes:
- JWT auth on all endpoints; namespace-scoped permissions.
dry-runevaluates but never fires actions β for testing rules pre-activation.- Event payloads encrypted at rest; PII masked in audit logs by default.
High-Level Design
FR-1: Author and Manage Rules
New components:
- Rules API β HTTP service for CRUD. Handles expression parsing at write time, versioning, validation.
- Expression Parser (ANTLR) β converts condition strings into validated ASTs. Write-time only.
- Postgres (rules + rule_versions) β durable store. Every edit creates a new immutable version.
flowchart LR
AUTHOR["Rule Author"]:::client
API["Rules API"]:::service
PARSER["Expression Parser<br/>ANTLR"]:::service
DB[("Postgres<br/>rules + versions")]:::data
AUTHOR -->|"1. POST rule with condition string"| API
API -->|"2. Parse condition to AST"| PARSER
PARSER -->|"3. Return validated AST"| API
API -->|"4. Store rule + version"| DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step:
- Author submits a rule with a condition string (e.g.,
order.total > 500 AND user.tier == 'NEW'). - Rules API sends condition to ANTLR parser. If syntax errors, returns 400 immediately.
- API validates AST (field references plausible, action configs valid), stores as new
rule_versionsrow. - Status starts DRAFT. Author calls
POST /activatewhen ready β flips to ACTIVE and publishesrule.updatedto Kafka.
Why version on every edit: Compliance requires knowing which logic was active at any point. Immutable versions make rollback trivial β just repoint current_version.
FR-2: Evaluate Events Against Rules
New components:
- Evaluation Engine β hot-path service. Loads compiled Rete network from memory, walks event through it. π‘ Rete = network of condition nodes where shared sub-conditions are evaluated once and results propagated.
- In-Memory Rule Cache β each engine holds the full compiled rule set. Rebuilt on startup or Kafka signal.
- Action Bus (Kafka) β matched actions published here, decoupling evaluation from execution.
flowchart LR
CALLER["Calling Service"]:::client
ENGINE["Evaluation Engine<br/>Rete network"]:::service
MEMCACHE[("In-Memory Cache<br/>compiled rules")]:::data
KAFKA["Action Bus<br/>Kafka"]:::async
AUDIT[("Postgres<br/>rule_executions")]:::data
CALLER -->|"1. POST /evaluate"| ENGINE
ENGINE -->|"2. Load Rete network"| MEMCACHE
ENGINE -->|"3. Write audit record"| AUDIT
ENGINE -->|"4. Publish matched actions"| KAFKA
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step:
- Calling service sends
POST /evaluatewith event payload. - Engine loads Rete network from local memory (zero network hop).
- Event fields asserted into alpha nodes. Shared conditions evaluated once.
- Beta nodes join results. Terminal nodes = fully-matched rules.
- Priority ordering applied. Returns result synchronously (< 50ms P99).
- Async: publishes matched actions to Kafka, writes audit record to Postgres.
Why in-memory, not Redis for the hot path: Each evaluation touches the full Rete network. At 0.5ms per Redis GET, 10K fetches = 5 seconds. The network must live in-process.
FR-3: Hot Reload Without Downtime
New components:
- Change Propagation (Kafka topic:
rule-changes) β Rules API publishes on every rule change. All engine instances consume. - Cache Rebuild Worker β background thread in each engine listens for changes, fetches updated set from Redis or Postgres, rebuilds Rete, atomically swaps reference.
flowchart LR
API["Rules API"]:::service
DB[("Postgres")]:::data
KAFKA["Kafka<br/>rule-changes"]:::async
REDIS[("Redis L2")]:::data
ENGINE1["Engine 1"]:::service
ENGINE2["Engine 2"]:::service
API -->|"1. Rule updated"| DB
API -->|"2. Publish change"| KAFKA
KAFKA -->|"3. Consume"| ENGINE1
KAFKA -->|"3. Consume"| ENGINE2
ENGINE1 -->|"4. Fetch set"| REDIS
ENGINE2 -->|"4. Fetch set"| REDIS
API -->|"5. Update L2"| REDIS
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step:
- Rules API updates Postgres and publishes
{ namespace, rule_id, action: "UPDATED", version: 7 }torule-changesKafka topic. - API serializes the updated rule set for that namespace to Redis (engines donβt all hammer Postgres).
- Each engine instance consumes the topic. On event: βis my in-memory version older?β
- If stale, fetch serialized rule set from Redis, deserialize, compile new Rete network.
- Atomically swap the in-memory reference (
AtomicReference.set()in Java, pointer swap in Go). In-flight evaluations finish on old network; new ones use new. - Total propagation: Kafka (~200ms) + Redis fetch (~2ms) + Rete compile (~50ms) β under 500ms.
Why Kafka, not Redis Pub/Sub: Kafka guarantees delivery even if an engine was briefly down (catches up on restart). Redis Pub/Sub is fire-and-forget β a missed message means stale rules until next change.
Technology Choices
| Tier / Purpose | What it stores | Primary pick | Alternatives |
|---|---|---|---|
| Rule definitions | rules, versions, AST as JSONB | PostgreSQL | MySQL, CockroachDB |
| Execution audit | which rules fired per event | PostgreSQL partitioned by day | Cassandra, ClickHouse |
| In-memory cache | compiled Rete network per namespace | Local process memory | Hazelcast |
| L2 cache | serialized rule set per namespace | Redis | Memcached |
| Change propagation | rule.updated events | Kafka | Redis Pub/Sub, NATS |
| Action dispatch | matched actions to execute | Kafka | SQS, RabbitMQ |
| Expression parsing | tokenize conditions at write time | ANTLR | Custom recursive descent |
| Dead letter queue | failed actions after max retries | Kafka DLQ topic | SQS DLQ |
Why pre-parsed AST in Postgres, not raw expression strings?
Parsing at evaluation time is expensive and error-prone. By parsing at write time and storing the AST as JSONB, the engine never touches a parser on the hot path. It loads a tree, walks it. Authors get syntax errors immediately, not when a customer triggers the rule.
Why local in-memory cache, not just Redis?
Redis ~0.5ms per GET. With 10K rules and a Rete network, reconstructing from serialized data per request is too slow. Each engine holds the full compiled network in memory. Redis is L2 for cold starts and consistency signals.
Data Modeling
Postgres Schemas
CREATE TABLE rules (
rule_id UUID PRIMARY KEY,
namespace VARCHAR(100) NOT NULL,
name VARCHAR(200) NOT NULL,
priority INTEGER DEFAULT 0,
status VARCHAR(15) NOT NULL DEFAULT 'DRAFT',
current_version INTEGER NOT NULL DEFAULT 1,
created_by VARCHAR(100) NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
updated_at TIMESTAMP NOT NULL DEFAULT NOW()
);
CREATE INDEX idx_rules_ns ON rules(namespace, status);
CREATE TABLE rule_versions (
version_id UUID PRIMARY KEY,
rule_id UUID NOT NULL REFERENCES rules(rule_id),
version_number INTEGER NOT NULL,
condition_ast JSONB NOT NULL, -- pre-parsed AST
actions JSONB NOT NULL,
created_by VARCHAR(100) NOT NULL,
created_at TIMESTAMP NOT NULL DEFAULT NOW(),
UNIQUE(rule_id, version_number)
);
CREATE TABLE rule_executions (
execution_id UUID PRIMARY KEY,
event_id VARCHAR(200) NOT NULL,
namespace VARCHAR(100) NOT NULL,
rules_evaluated INTEGER NOT NULL,
rules_matched INTEGER NOT NULL,
evaluation_time_ms REAL NOT NULL,
matched_rule_ids UUID[],
actions_triggered JSONB NOT NULL,
evaluated_at TIMESTAMP NOT NULL DEFAULT NOW()
) PARTITION BY RANGE (evaluated_at);
AST JSON Structure (stored in condition_ast JSONB)
{
"type": "AND",
"children": [
{
"type": "COMPARISON", "operator": ">",
"left": { "type": "FIELD_REF", "path": "order.total_amount" },
"right": { "type": "LITERAL", "value": 500, "dataType": "NUMBER" }
},
{
"type": "COMPARISON", "operator": "==",
"left": { "type": "FIELD_REF", "path": "user.tier" },
"right": { "type": "LITERAL", "value": "NEW", "dataType": "STRING" }
}
]
}
Supported node types: COMPARISON (>, <, >=, <=, ==, !=, IN, CONTAINS), ARITHMETIC (+, -, *, /), BOOLEAN (AND, OR, NOT), FIELD_REF (dot-path into event payload), LITERAL (typed constant), FUNCTION (built-ins like NOW(), LENGTH(), UPPER()).
Redis Cache Patterns
"rules:compiled:{namespace}" β Serialized rule set (protobuf). No TTL, invalidated via Kafka.
"rules:version:{namespace}" β Integer version counter. Used for cache-busting.
"rules:lock:{namespace}" β Distributed lock during reload (TTL: 5s).
Deep Dives
Deep Dive 1: Fast Evaluation at Scale
Evaluate an event against 10K rules in < 50ms.
Bad β Linear Loop: Iterate every rule, evaluate condition tree. O(N) with N=10K rules Γ 5 conditions = 50K checks. ~200ms. Fails P99 budget.
Good β Compiled AST with Short-Circuit:
Pre-compile each AST into a function. Evaluate in parallel (thread pool), short-circuit AND/OR early. Gets to ~30ms with 8 threads. Problem: many rules share sub-conditions (e.g., 200 rules all check user.country IN high_risk_list). Same condition evaluated 200 times.
Great β Rete Algorithm: Build a discrimination network at compile time:
- Alpha nodes β one per unique condition.
user.country IN high_risk_listexists once regardless of how many rules use it. - Beta nodes β join multiple alpha results. Fire when ALL required alphas match.
- Terminal nodes β fully-matched rules. Collect actions.
Event fields asserted once into the network. Shared conditions evaluated once, result propagates to all dependents. For 10K rules with 60% overlap, effective evaluations drop from 50K to ~8K. P99: 12-20ms.
Trade-off: Rete networks use more memory (~50MB per namespace for 10K rules). Acceptable for in-memory engines.
Deep Dive 2: Hot Reload Without Downtime
How to update rules across a fleet without restart or stale evaluations.
Bad β Restart on Change: Deploy new rules as config, rolling-restart all engine pods. Takes 5-10 minutes for 50 pods. During rollout, some pods have old rules, some new. For fraud rules, a 10-minute gap means money lost.
Good β Polling with Version Check: Each engine polls Redis every 5s: βhas namespace version incremented?β If yes, fetch and rebuild. Problem: 5-second staleness. Under batch releases (20 rules at once), all instances fetch simultaneously β thundering herd.
Great β Event-Driven Cache Swap (as in FR-3): Key refinements beyond the basic Kafka flow:
- Debounce: 20 rules change in 2 seconds? Wait 500ms after last event, one rebuild not twenty.
- Atomic swap: New Rete network built on background thread. Only swap when fully compiled. In-flight evaluations never interrupted.
- Distributed lock: Redis lock (
rules:lock:{namespace}, 5s TTL) prevents multiple engines from fetching from Postgres if Redis is cold. - Blocking rule consistency: For BLOCK_TRANSACTION rules, engine can optionally version-check against Postgres before returning a PASS. Adds ~5ms but ensures no fraud rule is missed.
Deep Dive 3: Expression Parsing
Turn order.total > 500 AND user.tier == 'NEW' into a validated, executable AST.
Bad β String eval() at Runtime:
Store condition as raw string, eval() at evaluation time. Security nightmare (injection), terrible performance (re-parse every call), impossible to validate field references statically.
Good β Custom Recursive Descent Parser: Hand-rolled tokenizer β recursive descent β AST. Works, but inconsistent operator precedence, poor error messages, and every new operator requires manual parser changes.
Great β ANTLR Grammar with Visitor Pattern: Define a formal grammar in ANTLR (expression β comparison | boolean combinator | parenthesized sub-expression). ANTLR generates lexer + parser. A custom Visitor walks the parse tree and builds our AST JSON structure. Benefits:
- Unambiguous precedence β defined in grammar, not hacked in code.
- Rich error messages β ANTLR pinpoints the exact failing token. Authors see βunexpected βANNDβ at position 15, did you mean βANDβ?β
- Extensible β new function (
LENGTH(field) > 5) = one grammar rule + one visitor method. - Semantic validation β after parsing, check field references against known event schemas, verify type compatibility.
Parsing runs once at write time (~2ms per rule). Stored AST loaded at evaluation time with zero parsing overhead.
Deep Dive 4: Action Execution Reliability
When a rule matches, trigger actions reliably without blocking evaluation.
Bad β Inline Execution: Fire webhook synchronously before returning response. Webhook timeout (3s) destroys evaluation P99. Action failure shouldnβt block the evaluation.
Good β Async Fire-and-Forget: Publish to Kafka, return immediately. If action fails, itβs lost. For fraud alerts thatβs unacceptable.
Great β Async with Retry + DLQ + Idempotency:
- Engine publishes matched actions to Kafka
actionstopic withaction_id(idempotency key),rule_id,event_id,action_config. - Action Executor consumes and executes (webhook, email, etc.).
- On failure: republish to
actions-retrywith exponential backoff (1s β 5s β 30s β 5min). Max 5 attempts. - After max retries: publish to
actions-dlq. Alert rule owner. - Idempotency: executor checks
action_idin Redis set before executing. If present (duplicate), skip. TTL: 24h. - For BLOCK_TRANSACTION specifically: decision returned synchronously in evaluation response. Kafka handles only notification/alerting.
flowchart LR
ENGINE["Engine"]:::service
KAFKA["actions topic"]:::async
EXECUTOR["Action Executor"]:::service
RETRY["actions-retry"]:::async
DLQ["actions-dlq"]:::async
TARGET["External Target"]:::client
ENGINE -->|"publish"| KAFKA
KAFKA -->|"consume"| EXECUTOR
EXECUTOR -->|"execute"| TARGET
EXECUTOR -.->|"failure"| RETRY
RETRY -->|"re-deliver"| EXECUTOR
EXECUTOR -.->|"max retries"| DLQ
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#AB47BC,stroke:#4A148C,color:#fff
Core Flows
Flow 1: Create and Activate a Rule
sequenceDiagram
actor Author
participant API as Rules API
participant Parser as ANTLR Parser
participant PG as Postgres
participant Kafka
participant Redis
participant Engine
Author->>API: POST /v1/rules (condition + actions)
API->>Parser: Parse condition string
Parser-->>API: Validated AST
API->>PG: INSERT rules + rule_versions (DRAFT)
API-->>Author: 201 Created
Author->>API: POST /v1/rules/:id/activate
API->>PG: UPDATE status=ACTIVE
API->>Kafka: Publish rule.activated
API->>Redis: Store updated rule set
API-->>Author: 200 OK
Kafka->>Engine: Consume rule.activated
Engine->>Redis: Fetch updated set
Engine->>Engine: Compile Rete + atomic swap
Walkthrough:
- Author submits a rule with a human-readable condition. API parses via ANTLR, validates, stores as DRAFT.
- Author tests via
POST /evaluate/dry-runto confirm behavior. - Author activates. API updates Postgres, publishes to Kafka, pre-populates Redis.
- Engine instances consume the event, fetch from Redis, rebuild Rete network, swap atomically.
- Next evaluation uses the new rule. Total time from activation to live: < 2 seconds.
Failure path: If Kafka is unavailable when API publishes the change event, the rule is active in Postgres but engines donβt know. Safety net: engines do a periodic full-sync every 60s (poll Postgres namespace version), bounding staleness even if Kafka has an outage.
Flow 2: Evaluate an Event
sequenceDiagram
actor Caller
participant Engine
participant Cache as In-Memory Cache
participant Kafka as Action Bus
participant PG as Postgres
participant Executor as Action Executor
participant Target as Webhook
Caller->>Engine: POST /evaluate (namespace + payload)
Engine->>Cache: Load Rete network
Note over Engine: Assert fields into alpha nodes
Note over Engine: Propagate through beta to terminals
Engine-->>Caller: 200 matched rules + decision (12ms)
par Async
Engine->>Kafka: Publish matched actions
Engine->>PG: Write audit record
end
Kafka->>Executor: Consume action
Executor->>Target: POST webhook
alt Webhook fails
Executor->>Kafka: Republish to retry topic
end
Walkthrough:
- Calling service sends event. Engine loads Rete network from local memory.
- Event fields asserted into alpha nodes. Shared conditions evaluated once.
- Matched rules collected, priority-sorted. Synchronous response in ~12ms.
- Async: actions published to Kafka, audit record written to Postgres.
- Action Executor delivers webhooks with retry logic.
Failure path: Engine crashes after returning evaluation but before publishing actions to Kafka. Caller got BLOCK (correct), but fraud team webhook never fires. Mitigation: write an outbox record to Postgres alongside the audit log. Outbox poller picks up un-published actions and retries.
Interview Tips
-
Lead with the AST insight. βIβd store rules as pre-parsed ASTs, not raw strings. Parsing at write time, evaluation is just tree traversal.β Shows hot-path vs cold-path split understanding.
-
Name Rete. Most candidates loop linearly. Mention shared sub-condition optimization. Amazon and Goldman specifically look for this.
-
Draw the Kafka boundary early. Evaluation is synchronous and fast; action execution is async and reliable. Shows decoupling principle.
-
Hot reload is the trap question. Have the Kafka + atomic swap + version check story ready.
-
Dry-run is your testing story. Prevents bad rules from breaking production. Version history + instant rollback + AST validation at write time.
-
Compliance at financial companies. At Goldman, PayPal, Razorpay: immutable audit trail. Every evaluation logged. Every version preserved. Answer βwhat rule was active at 3:47 PM March 12?β from
rule_versions. -
Anchor with scale numbers. 10K active rules, 50K evals/sec, < 50ms P99. Realistic for a payment fraud engine.
Discussion
Newest first