← Projects

BARO — Autonomous Taxi Dispatch & Control Platform

  • Apr – Jun 2026
  • Team of 5 · Backend / DevOps
  • AWS ECS · Terraform
  • MQTT · Kafka
  • OpenStack · K3s

Overview

For autonomous taxis to work on real roads, you have to receive the positions of hundreds of vehicles without gaps, dispatch the nearest idle car, and send cars that finish a ride back to where demand is high. BARO ties that whole loop — request → dispatch → ride → relocation — into one driverless taxi dispatch and control platform.

Instead of real cars, a Python simulator runs 1,500 vehicles starting from 254 taxi stands in Seoul. The services run on AWS (ECS Fargate plus self-hosted Kafka and Mosquitto), while analytics, observability and long-term storage run on K3s on an on-prem OpenStack cloud. The two sides are connected by a site-to-site VPN.

It was the final cloud project of the 3rd cohort of the Hyundai AutoEver Mobility SW School, built by five people from April to June 2026. I worked on backend and DevOps.

Architecture

Scroll and the diagram on the right follows the text on the left. Hover over a node to see its specs.

At a glance

BARO splits into rider apps and vehicles, AWS where the services run, and on-prem (OpenStack K3s) for analytics and observability. AWS and on-prem are joined by a site-to-site VPN.

ECS Fargate ×7EC2: Kafka · MosquittoIPsec VPN

1How vehicle positions flow

1,500 vehicles publish their position over MQTT (QoS 0) every 3 seconds. control receives them through a shared subscription, so each message goes to only one instance even as it scales out.

vehicles/{id}/telemetry$share/control-service

The control console gets positions straight over SSE; everything else goes to the Kafka vehicle-data-topic. Keyed by carId, so each vehicle’s positions stay in order.

4 partitionskey = carId

dispatch drops messages older than 10 seconds and keeps only idle cars in the Valkey GEO index. The on-prem consumer reads the same topic across the VPN and stores it in TimescaleDB.

GEO dispatch:cars:idle:geohypertable vehicle_data

2From one request to a dispatched car

A request passes the ALB and then the gateway. JWTs are issued by user and verified only at the gateway. dispatch first shows a fare and route from Kakao (the pre-dispatch quote).

X-Authenticated-User-Idquote TTL 10 min

On confirmation it finds 10 candidates within 15 km, skips cars silent for over 300 seconds, and takes a SETNX lock (30 s) so two riders can’t grab the same car. The dispatch is saved in a transaction.

GEOSEARCH 15kmSETNX TTL 30s

The command reaches the car over MQTT (QoS 1) and an ACK comes back. If no ACK arrives within 10 seconds, the next car is dispatched automatically.

vehicles/{id}/commandsACK timeout 10s

The rider app gets live vehicle positions and arrival status (pickup → destination) over SSE.

SSE vehicle-location

3After the ride — relocation

When a car reports arriving at its destination, control notifies relocation asynchronously, kept separate from the dispatch flow.

ARRIVED(to_dest)202 Accepted

PostGIS scores nearby stands (0.7 · demand − 0.3 · distance) and the car is sent to the best one with RELOCATE. It stays dispatchable while moving, so it can take a new request right away.

ST_DWithin 10→30kmstatus = relocating

4Overnight — learning demand

Every day at 02:00, Airflow pulls the previous day’s dispatch data through the internal ALB across the VPN and aggregates demand by stand × weekday × time slot.

gzip CSV exportX-Internal-Api-Key

Stand weights learned with XGBoost are sent to relocation and apply from the next relocation onward.

weights 0–1

5Keeping watch — observability

On-prem Prometheus scrapes service, Kafka and EC2 metrics across the VPN, and alerts go to three Slack channels: AWS, on-prem and services.

JMX :9404CloudWatch exporterAlertmanager → Slack
At a glance1 / 13
Clients · vehiclesAWS · ap-northeast-2VPC 10.20.0.0/16On-prem · K3sSITE-TO-SITE VPN · IPSECRider web appReact · SSE clientControl consolebaro-admin · SSEVehicles ×1500asyncio simulatorKakao MobilityRoute & fare APISlack3 alert channelsPublic ALBHTTPS · path routinggatewayJWT · circuit breakeruserAccounts · tokensMosquittoEC2 · MQTTdispatchDispatch · lockscontrolMQTT ↔ Kafka hubValkeyElastiCache · GEORDS PostgresPostGIS · schemasKafkaKRaft · 4 partitionsrelocationStand scoringInternal ALB/internal · metricsPrometheusGrafana · alert rulesAirflowDaily 02:00 demand modelTimescaleDBTime-series hypertablekafka-consumerConsumes across VPN

Try it

A demo where hundreds of cars move on a map and requests, dispatch, rides and relocation play out on their own

Coming in the next update.

Operations log

When Kafka stalled, the control console stalled too

Symptom
The control console showed only one vehicle, and Mosquitto logs filled with broken pipes.
Cause
The Kafka producer’s default wait (max.block.ms, 60 s) held the MQTT receive thread, and the broker’s queue overflowed.
Fix
Cut max.block.ms to 500 ms and retries to 0 so a Kafka outage can no longer spill into MQTT receiving.
Result
With Kafka down, MQTT receiving and the console’s SSE keep working.

control-service fans each MQTT position out to two places: the control console (SSE) and Kafka. When the Kafka broker dropped, the producer waited for metadata and blocked in send() for up to 60 seconds — inside the MQTT receive thread. While receiving stalled, Mosquitto’s queue filled up and the connection broke (broken pipe), leaving the console with only the last vehicle it had heard from.

Kafka only stores positions, so losing a few briefly is acceptable; the console must never freeze. Publishing now fails fast when it takes too long, and failures are visible as a metric (baro_control_telemetry_kafka_publish_failed_total).

Lesson A default timeout on a synchronous call can carry a failure into the next system. I set timeouts per path instead of trusting defaults.

Consumer lag hit 2.9M — the culprit was a DB query on the hot path

Symptom
dispatch fell behind on positions and Kafka consumer lag grew to 2.9M.
Cause
For every position message (about 333 per second), it looked up the vehicle’s active dispatch in the database.
Fix
Cached the vehicle → dispatch mapping in memory so message handling no longer touches the database.
Result
The per-message query disappeared and the lag cleared.

dispatch reads every position from Kafka and, if the car is on a trip, pushes it to the rider app over SSE. The problem was that it asked the database “is this car on a trip?” for every message. With 1,500 cars reporting every few seconds, that meant hundreds of queries per second; consumption fell behind production and the lag ballooned to 2.9M.

Now the mapping is updated only when a dispatch starts or ends, and message handling reads memory only.

Lesson On a path that runs hundreds of times a second, even one query is expensive. Clear I/O off the hot path first.

At 2,000 vehicles, Mosquitto started dropping messages

Symptom
A load test with 2,000 vehicles showed dropped messages.
Cause
Messages beyond Mosquitto’s default max_queued_messages (1,000) were being discarded.
Fix
Raised the limit to 10,000 and then 50,000. Unlimited (0) was rejected because of the out-of-memory risk.

While raising the vehicle count to 2,000, the broker was trimming undelivered messages at its queue limit. Raising the limit is easy, but unlimited means broker memory grows without bound whenever a subscriber slows down. On a single t3.micro broker I capped it at 50,000, and staggered vehicle connections by 0.05 s per 10 cars to soften connection storms.

Lesson A queue limit is a "drop when full" policy. Even when raising it, pick a ceiling memory can actually hold instead of going unlimited.

Design decisions

Running Mosquitto on EC2 instead of AWS IoT Core

Context
Vehicles first talked to AWS IoT Core (mTLS, X.509 certificates). As the simulated fleet grew, connection rate limits forced a 1.2 s delay per 10 cars, and costs scaled with the number of vehicles.
Alternatives considered
  • Keep AWS IoT Core (the original setup)
Why
I control the broker’s settings (queue limits, auth, sessions), and without a connection rate limit 1,000 cars connect in about 5 seconds at 0.05 s per 10 cars.
Trade-off
Availability, patching and monitoring are on us. Credentials live in Secrets Manager and deployment is automated through SSM to keep the load down.

Moving Kafka from ECS Fargate to EC2 + EBS (single-node KRaft)

Context
Kafka first ran as an ECS Fargate task with data on EFS. The advertised listener address clients connect to was hard to pin down, and the log storage sat on a network file system.
Alternatives considered
  • Keep ECS Fargate + EFS (the original setup)
  • Amazon MSK
Why
A fixed private IP and a Cloud Map name (kafka.baro.internal) pin the advertised listener, and logs sit on EBS gp3. vehicle-data-topic has four partitions to match dispatch’s consumer concurrency of four.
Trade-off
With one broker (RF=1), a broker failure stops the position stream. Positions are ephemeral and kept for only an hour, so cost came first.

QoS 0 for positions, QoS 1 for commands — plus shared subscriptions

Context
Positions pour in every few seconds from every car and are soon replaced by newer ones, while a single lost dispatch command, arrival event or ACK breaks a dispatch. The control service also had to scale to several instances.
Alternatives considered
  • QoS 1 for every message
  • Every control instance subscribing to all topics
Why
Positions go light at QoS 0; commands, events and ACKs go reliably at QoS 1. control uses a shared subscription ($share/control-service/…) and a per-instance clientId, so with several instances each message is handled by only one.
Trade-off
QoS 0 positions are lost while a connection is down. Cars keep the last 200 in a buffer and resend them on reconnect over a QoS 1 topic (telemetry/buffered).

SSE instead of WebSocket for live updates

Context
The rider app and the control console only need to receive positions and status from the server.
Alternatives considered
  • WebSocket
Why
SSE is enough for one-way push, and because it is plain HTTP it passes through nginx and the Vercel function proxy unchanged. Clients reconnect with exponential backoff when dropped.
Trade-off
SSE connections live in server memory, so scaling out splits them across instances (see Limits & retrospective).

Authenticate once at the gateway, block internal paths twice

Context
With five services, verifying JWTs in each one would scatter verification logic and key management.
Alternatives considered
  • Verify JWTs in every service
Why
JWTs are verified only at the gateway, and user details travel in X-Authenticated-User-* headers; if a client sends headers with those names, the gateway strips them. Service-to-service calls are checked with X-Internal-Api-Key, and internal paths are blocked twice — at the ALB (403) and the gateway (404).
Trade-off
Services trust the headers on the assumption that they sit behind the gateway, so ALB and security group rules must keep internal paths from ever being exposed.

One RDS instance, one schema per service

Context
A database instance per service multiplies cost by the number of services.
Alternatives considered
  • An RDS instance per service
Why
One RDS PostgreSQL instance (db.t4g.micro) holds separate user, dispatch, relocation and control schemas so services don’t touch each other’s tables. A one-off ECS task (db-init) creates the schemas.
Trade-off
Every service shares database failures and load. Splitting into per-service instances is the next step as traffic grows.

Infrastructure & cost

git pushmainActionsGitHub · path filterBuild · testJDK 21 · GradleECRimage per serviceECS deploypassing CI only

AWS, built with Terraform

  • A single envs/dev environment is managed in Terraform: the VPC (2 AZs, public and private subnets), public and internal ALBs, seven ECS Fargate services, RDS PostgreSQL, ElastiCache (Valkey), EC2 for Kafka and Mosquitto, Cloud Map, Secrets Manager and the site-to-site VPN.
  • State lives in S3 with a DynamoDB lock. GitHub Actions runs fmt · validate · plan on pull requests and apply → db-init → edge deploy (SSM) on main.
  • To prevent accidents, destroy only runs after typing a confirmation string such as destroy-dev.

Deployment pipeline

The server pipeline detects changed paths, builds and tests only those services, and ships only the ones that pass CI to ECS through ECR. The rider web app deploys with CodeDeploy blue/green and rolls back automatically on failure.

Keeping costs down

  • One switch, runtime_enabled=false, tears down NAT, the ALBs, ECS, RDS and EC2 while keeping what is tedious to recreate — the VPC, ECR and secrets.
  • An RDS snapshot is taken automatically before teardown, and the next apply restores the latest one.
  • Kafka’s EC2 went from t3.medium to t3.small, the bastion was removed in favor of SSM-only access, log retention was cut to 7 days, and per-vehicle logs moved to DEBUG to lower CloudWatch costs.

My contribution

This section is coming soon.

Limits & retrospective

Known limitations, and what I would change next time.

  • In-memory state blocks horizontal scaling. control’s vehicle state, dispatch’s pending-ACK list and SSE connections all live in JVM memory, so adding instances splits the state. Moving them to a shared store such as Redis or pub/sub is the next step.
  • Single points of failure remain. Kafka runs on one broker (RF=1), there is one NAT gateway, and RDS is single-AZ. That was a cost-first choice; production would need three brokers, a NAT per AZ and Multi-AZ RDS.