Gig-Queue is a distributed ticket pre-sale system built on Apache Kafka. It simulates a flash sale ticket scenario: thousands of buyers competing for a limited number of seats in a few seconds.
Kafka is used not as a transport but as an ordered, persistent registry of the purchase queue: the offset the broker assigns to a request is its position in line. Fairness is therefore a property of the log, not something the application reconstructs.
For more dettails please check the report.
| Service | Role |
|---|---|
ticket-api |
FastAPI edge. POST /buy → 202, GET /status → queue position and outcome |
stress-producer |
synthetic load generator |
inventory-service (×3) |
allocates seats through an atomic Lua script |
fraud-detector |
fixed-window rate counters per user and per IP |
notifier |
consumes outcomes and alerts, sends email through SMTP |
dlq-monitor |
accounts for discarded payloads |
dashboard |
read-only console: cluster health and aggregated state |
| Endpoint | Address |
|---|---|
| API | http://localhost:8080 |
| Console | http://localhost:8081 |
| Mailpit | http://localhost:8025 |
| AKHQ | http://localhost:8090 |
| Topic | Partitions | RF | min.ISR | Key | Rationale |
|---|---|---|---|---|---|
topic-requests |
3 | 3 | 2 | event_id |
total ordering per event; parallelism enabled across events |
topic-orders |
3 | 3 | 2 | order_id |
computed order synchronously in the inventory loop, so its leadership must not concentrate on one broker |
topic-fraud |
1 | 3 | 2 | user_id |
separate alert stream, not a mandatory service |
topic-dlq |
1 | 3 | 2 | — | payloads that cannot be deserialized |
- Docker, Docker Compose
openssl,keytool- Python 3 with
requests,redisandconfluent-kafkafor the tests
The order matters: the brokers must be healthy before the topics are created, the topics before the ACLs, and the ACLs before the services connect. A service that starts without its permissions fails authorization and restarts in a loop.
# 1. the keystore password, read by generate-certs.sh and by docker compose
echo 'KAFKA_SSL_PASSWORD=gigQueueSecret' > .env
# 2. private CA, broker certificates with SAN, one client certificate per service
cd security && ./gen-certs.sh -force && cd ..
# 3. the cluster on its own; wait until all three report healthy
docker compose up -d --build kafka-1 kafka-2 kafka-3 redis mailpit
docker compose ps
# 4. topics and partition leadership (runs in the foreground, exits when done)
docker compose up kafka-init
# 5. the authorization matrix, one row per principal
bash ./security/create-acls.sh
# 6. the services, now that their permissions exist
docker compose up -d --buildcurl -s localhost:8080/healthz
docker exec kafka-1 kafka-topics --bootstrap-server kafka-1:9093 \
--command-config /etc/kafka/secrets/admin.properties \
--describe --topic topic-requestsThree partitions, three in-sync replicas each, and leaders on different brokers. Then open the console at http://localhost:8081.
docker compose up -d is enough for a normal restart. But after
docker compose down -v or generate-certs.sh --force, restart every
service: a producer still holding an idempotence session against a cluster
whose logs were wiped enters a fatal state and never recovers, and every
POST /buy fails until the process is replaced.
docker compose restart ticket-api inventory-service fraud-detector \
notifier dlq-monitor dashboardThree independent layers.
-
Encryption: TLS on every listener, including the KRaft controller traffic. A plaintext client is disconnected during the handshake on all three external listeners.
-
Authentication.
ssl.client.auth=requiredwith a private CA. Each service holds its own certificate;ssl.principal.mapping.rulesextracts the CN, so the Kafka principal is the service name. A syntactically valid certificate signed by an unknown CA is rejected during the handshake even when its CN impersonates a real service. -
Authorization.
StandardAuthorizerwith deny-by-default and one row per principal. For example, thedashboardholdsDescribeand neverRead: it reports how many replicas are in sync without being able to read a single message.
Each container mounts only its own certificate and private key, not the
security/ directory, so a compromised service cannot read anyone else's
credentials.
Each property is enforced by a configuration and verified by a demo, not asserted.
| Property | How it is enforced |
|---|---|
| High availability | RF=3 across three brokers, automatic leader election among in-sync replicas; three inventory replicas in one consumer group |
| Durability | acks=all + min.insync.replicas=2: an acknowledged record exists on at least two independent replicas, and writes are refused when the cluster cannot replicate them |
| Consistency | seat allocation runs inside a single atomic Lua script; processed:{order_id} stores the allocated range, so a redelivery returns it instead of decrementing again |
| Fairness (FCFS) | event_id as the message key puts every request for an event on one partition, where offsets are a total order |
| Scalability | partitioning by event; the consumer group scales up to the partition count, and service state lives in Redis so instances are interchangeable |
| Security | mTLS on every listener, one certificate per service, deny-by-default ACLs |
| Observability | services maintain their own aggregates in Redis; the console composes them with the Kafka Admin API without recomputing anything |
cd services/stress-producer
# fairness and head-of-line blocking
python3 stress-producer.py flash-sale
# rate limiting at the edge
python3 stress-producer.py bot-attack
# high availability
python3 stress-producer.py broker-failure
# deliberate redelivery
python3 stress-producer.py replay-idempotence
# encryption and authentication
python3 ../../scripts/demo_mtls.py
# authorization
python3 ../../scripts/demo_acl.py Every scenario ends with three automated checks: no acknowledged order is left without an outcome, seat ranges tile 1..N with no gap or overlap, and seats are assigned in the order the broker recorded the requests.
flash-sale picks four events so that two land on the same partition, and
gives three of them identical light traffic. The one sharing a partition with
the rush is delayed; the two on other partitions are not.
cd tests
python3 test_ticket_api.py
python3 test_inventory.py
python3 test_notifier.py
python3 test_dlq.py
python3 test_ssl.py
python3 test_fraud.py # last: it leaves 5-minute blocks behindPlain functions and assertions with a small runner (no framework).
