17 KiB
100K Vehicle Ingest Hardening Implementation Plan
For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: Move the Go vehicle ingest platform toward stable production operation for 100,000 vehicles.
Architecture: Keep the current Gateway -> NATS -> Kafka -> TDengine/Redis/MySQL architecture, but make capacity explicit, add reproducible load testing, harden OS/systemd/network limits, and remove single-threaded or per-message bottlenecks from storage projectors. The first target is predictable behavior under high connection count and high frame rate; horizontal scaling comes after measurable single-node limits.
Tech Stack: Go 1.26, systemd, Linux TCP tuning, NATS JetStream, Kafka, Redis, TDengine, MySQL, Prometheus-style /metrics.
Current Baseline
Observed on 2026-07-03 from ECS 115.29.187.205:
| Item | Current Observation |
|---|---|
| Gateway active connections | JT808 about 242, GB32960 about 2 |
| ECS load | about 0.09 / 0.15 / 0.21 |
| Memory | about 7.5 GB total, 6.6 GB available |
| Disk | / about 37% used |
| NATS bridge ack pending | 0 |
| NATS bridge pending | about 26 |
| Kafka lag | 0 across history/realtime/stat sampled metrics |
| Redis projector latency | about 0 ms gauge |
| MySQL projector latency | about 1-2 ms gauge |
| Gateway file limit | 1048576 |
| Gateway max processes | 30123 |
Kernel net.core.somaxconn |
128 |
| Gateway listen backlog | Go net.Listen default, bounded by somaxconn |
| Gateway default max connections | env default adjusted to TCP_MAX_CONNECTIONS=120000 |
Key gap: production is healthy at current load, but it has not proved 100K connection or high burst behavior. The current OS accept backlog and gateway default connection cap are below the target.
Capacity Assumptions
Use these as starting targets. Revise after load test evidence.
| Scenario | Target |
|---|---|
| Connected vehicles | 100,000 total |
| Average frame interval | 10-30 seconds depending protocol |
| Average ingress FPS | 3,000-10,000 frames/s |
| Burst ingress FPS | 20,000 frames/s for short periods |
| Gateway p99 parse + enqueue | < 50 ms |
| Gateway publish queue drops | 0 under target load |
| NATS ack pending | sustained 0, short spikes acceptable |
| Kafka lag | bounded and decreasing after burst |
| Redis current-state lag | < 3 seconds under normal load |
| TDengine history lag | bounded and observable |
| MySQL snapshot lag | can be lower priority and async; must not block Redis |
File Map
| File | Responsibility |
|---|---|
docs/superpowers/plans/2026-07-03-100k-vehicle-ingest-hardening.md |
This implementation plan |
docs/ops/go-vehicle-ingest-memory.md |
Operational memory, update after each production change |
docs/ops/vehicle-ingest-runbook.md |
Production runbook, add capacity and tuning commands |
go/vehicle-gateway/cmd/load-sim/main.go |
New load simulator for JT808/GB32960 TCP and MQTT-like payload replay |
go/vehicle-gateway/internal/loadsim/* |
New reusable load simulation helpers |
go/vehicle-gateway/internal/gateway/tcp_server.go |
Gateway connection handling, metrics, connection admission |
go/vehicle-gateway/cmd/gateway/main.go |
Gateway env defaults for 100K mode |
go/vehicle-gateway/internal/eventbus/async_sink.go |
Async publish queue metrics/backpressure behavior |
go/vehicle-gateway/cmd/nats-fast-writer/main.go |
Fast writer batching and worker behavior |
go/vehicle-gateway/cmd/history-writer/main.go |
Kafka/TDengine consume/write batching |
go/vehicle-gateway/internal/history/writer.go |
TDengine write path; later batch insert |
go/vehicle-gateway/cmd/realtime-api/main.go |
Redis-first and MySQL async projector behavior |
deploy/systemd/*.conf or docs/ops/*.md |
Documented sysctl/systemd tuning until deployment templates exist |
Task 1: Capacity Baseline and Production Guardrails
Files:
-
Modify:
docs/ops/go-vehicle-ingest-memory.md -
Modify:
docs/ops/vehicle-ingest-runbook.md -
Create:
docs/ops/100k-capacity-baseline.md -
Step 1: Create the capacity baseline document
Create docs/ops/100k-capacity-baseline.md with:
# 100K 车辆接入容量基线
更新时间:2026-07-03
## 目标
支撑 100,000 台车辆接入,入口服务在高连接数和高帧率下保持可观测、可背压、可恢复。
## 初始容量假设
| 项 | 目标 |
| --- | --- |
| 连接车辆数 | 100,000 |
| 平均帧间隔 | 10-30 秒 |
| 平均入口 FPS | 3,000-10,000 |
| 短时 burst FPS | 20,000 |
| Gateway p99 parse + enqueue | < 50ms |
| Redis 当前态延迟 | < 3s |
| TDengine 历史写入 lag | 可观测且 burst 后下降 |
| MySQL 当前态 | 异步,不能阻塞 Redis |
## 当前生产观测
当前 ECS `115.29.187.205` 在低负载下健康:
- JT808 连接约 242。
- GB32960 连接约 2。
- Kafka lag 为 0。
- NATS ack pending 为 0。
- Redis 写入约 0ms。
- MySQL 投影约 1-2ms。
## 当前已知缺口
- `net.core.somaxconn=128`,无法作为 100K 连接生产基线。
- Gateway 默认 `TCP_MAX_CONNECTIONS` 已调整为 `120000`,生产仍可通过环境变量覆盖。
- 已新增可重复的 TCP 连接压测工具,后续需要用它跑 10K/50K/100K 阶段测试并记录结果。
- 仍缺按协议维度的连接生命周期、读超时、parse/publish p95/p99 直方图。
- TDengine history writer 当前逐消息插入,后续需要批量写入。
- Kafka topic 当前 12 分区,100K 目标下需要结合实际 FPS 再评估分区数。
- Step 2: Update runbook with capacity baseline link
Add this line near the top of docs/ops/vehicle-ingest-runbook.md after the production inventory link:
10W 车辆容量目标、当前缺口和压测口径见 [100K 车辆接入容量基线](100k-capacity-baseline.md)。
- Step 3: Verify docs render
Run:
rg -n "100K|10W|容量基线|somaxconn|TCP_MAX_CONNECTIONS" docs/ops docs/architecture
Expected: the new baseline and runbook link are found.
- Step 4: Commit
git add docs/ops/100k-capacity-baseline.md docs/ops/vehicle-ingest-runbook.md docs/ops/go-vehicle-ingest-memory.md
git commit -m "docs: add 100k ingest capacity baseline"
Task 2: Load Simulator Skeleton
Files:
-
Create:
go/vehicle-gateway/internal/loadsim/config.go -
Create:
go/vehicle-gateway/internal/loadsim/config_test.go -
Create:
go/vehicle-gateway/internal/loadsim/templates.go -
Create:
go/vehicle-gateway/internal/loadsim/template_test.go -
Create:
go/vehicle-gateway/internal/loadsim/schedule.go -
Create:
go/vehicle-gateway/internal/loadsim/schedule_test.go -
Create:
go/vehicle-gateway/internal/loadsim/runner.go -
Create:
go/vehicle-gateway/internal/loadsim/runner_test.go -
Create:
go/vehicle-gateway/cmd/load-sim/main.go -
Create:
go/vehicle-gateway/cmd/load-sim/main_test.go -
Step 1: Write failing tests for loadsim config, templates, schedule, runner and CLI stats
Implemented tests:
-
TestConfigFromFlagsParsesCapacityKnobs -
TestConfigFromFlagsRejectsUnsafeValues -
TestFrameTemplateReturnsParsableJT808LocationFrame -
TestFrameTemplateReturnsParsableGB32960RealtimeFrame -
TestConnectionBatchesSpreadsConnectionsByRate -
TestRunnerDialsTargetConnectionsAndWritesFrames -
TestFormatStatsIncludesCapacityCounters -
Step 2: Run tests to verify RED
Run:
cd go/vehicle-gateway
go test ./internal/loadsim
Expected: FAIL because package or symbols do not exist.
- Step 3: Implement minimal load simulator
Implemented:
-
flag config validation
-
parsable JT808 0x0200 and GB32960 realtime frame templates
-
connection batch scheduling
-
TCP runner with dial injection for tests
-
cmd/load-simCLI -
Step 4: Run test to verify GREEN
Run:
cd go/vehicle-gateway
go test ./internal/loadsim
Expected: PASS.
- Step 5: Add CLI skeleton
Implemented go/vehicle-gateway/cmd/load-sim/main.go.
- Step 6: Verify CLI builds
Run:
cd go/vehicle-gateway
go test ./cmd/load-sim ./internal/loadsim
Expected: PASS.
- Step 7: Commit
git add go/vehicle-gateway/internal/loadsim go/vehicle-gateway/cmd/load-sim
git commit -m "test(go): add vehicle ingest load simulator skeleton"
Task 3: Gateway Capacity Configuration
Files:
-
Modify:
go/vehicle-gateway/cmd/gateway/main.go -
Modify:
docs/ops/100k-capacity-baseline.md -
Modify:
docs/ops/vehicle-ingest-runbook.md -
Step 1: Write test for 100K connection env default
Added TestGatewayDefaultsTo100KConnectionCeiling.
- Step 2: Run test to verify RED before changing default
Run:
cd go/vehicle-gateway
go test ./cmd/gateway
Expected: FAIL while call site still used 20000.
- Step 3: Change gateway default
In go/vehicle-gateway/cmd/gateway/main.go, change the TCP server config default:
MaxConnections: envInt("TCP_MAX_CONNECTIONS", 120_000),
- Step 4: Document required sysctl
Add to docs/ops/100k-capacity-baseline.md:
## ECS OS 参数建议
100K 连接目标需要至少以下系统参数作为起点:
```bash
sysctl -w net.core.somaxconn=65535
sysctl -w net.ipv4.tcp_max_syn_backlog=65535
sysctl -w net.ipv4.ip_local_port_range="10000 65000"
sysctl -w net.ipv4.tcp_tw_reuse=1
systemd gateway 已配置 LimitNOFILE=1048576,需要持续保持。
- Step 5: Run tests
Run:
cd go/vehicle-gateway
go test ./cmd/gateway ./internal/gateway
Expected: PASS.
- Step 6: Commit
git add go/vehicle-gateway/cmd/gateway/main.go go/vehicle-gateway/cmd/gateway/main_test.go docs/ops/100k-capacity-baseline.md docs/ops/vehicle-ingest-runbook.md
git commit -m "chore(go): raise gateway connection target"
Task 4: Gateway Runtime Metrics for 100K Operations
Files:
-
Modify:
go/vehicle-gateway/internal/gateway/tcp_server.go -
Modify:
go/vehicle-gateway/internal/metrics/metrics.go -
Test:
go/vehicle-gateway/internal/gateway/tcp_server_test.go -
Test:
go/vehicle-gateway/internal/metrics/metrics_test.go -
Step 1: Add metrics test for rejected connections
Added TestTCPServerRecordsConnectionRejectionMetric in go/vehicle-gateway/internal/gateway/tcp_server_test.go.
- Step 2: Run RED or confirm helper coverage
Run:
cd go/vehicle-gateway
go test ./internal/gateway -run TestTCPServerRecordsConnectionRejectionMetric
Expected: FAIL before recordConnectionRejection exists.
- Step 3: Record max connection rejections
In go/vehicle-gateway/internal/gateway/tcp_server.go, when the semaphore default case rejects a connection, add:
server.recordConnectionRejection("max_connections")
Implemented recordConnectionRejection.
- Step 4: Run gateway tests
Run:
cd go/vehicle-gateway
go test ./internal/gateway ./internal/metrics
Expected: PASS.
- Step 5: Commit
git add go/vehicle-gateway/internal/gateway/tcp_server.go go/vehicle-gateway/internal/gateway/tcp_server_test.go go/vehicle-gateway/internal/metrics/metrics.go go/vehicle-gateway/internal/metrics/metrics_test.go
git commit -m "feat(go): expose gateway connection rejection metrics"
Task 5: Redis-First Realtime Projector Safety
Files:
-
Modify:
go/vehicle-gateway/cmd/realtime-api/main.go -
Test:
go/vehicle-gateway/cmd/realtime-api/main_test.go -
Step 1: Add test proving MySQL queue drops do not fail Redis update
Add a test around asyncSecondaryRealtimeUpdater in go/vehicle-gateway/cmd/realtime-api/main_test.go using existing test patterns. The test should:
- Create an updater with queue size
1. - Fill the queue.
- Call
Updatewith a new envelope. - Assert
Updatereturnsnil. - Assert metric
vehicle_realtime_async_queue_total{status="dropped"}increments.
- Step 2: Run test
Run:
cd go/vehicle-gateway
go test ./cmd/realtime-api -run TestAsyncSecondaryQueueDropDoesNotFailPrimaryUpdate
Actual: PASS. The behavior already existed; this test locks it as regression coverage.
- Step 3: Ensure queue drop is observable and non-fatal
In go/vehicle-gateway/cmd/realtime-api/main.go, preserve current Redis-first behavior:
-
Redis update remains synchronous and primary.
-
MySQL update remains async secondary.
-
Full secondary queue increments dropped metric and does not return error to Kafka processing.
-
Step 4: Run tests
cd go/vehicle-gateway
go test ./cmd/realtime-api ./internal/realtime
- Step 5: Commit
git add go/vehicle-gateway/cmd/realtime-api/main.go go/vehicle-gateway/cmd/realtime-api/main_test.go
git commit -m "test(go): lock redis-first realtime projector behavior"
Task 6: TDengine History Writer Batching Design
Files:
-
Create:
docs/architecture/tdengine-batch-writer-design.md -
Later Modify:
go/vehicle-gateway/cmd/history-writer/main.go -
Later Modify:
go/vehicle-gateway/internal/history/writer.go -
Step 1: Document the batch writer decision
Created docs/architecture/tdengine-batch-writer-design.md.
- Step 2: Commit design before code
git add docs/architecture/tdengine-batch-writer-design.md
git commit -m "docs: design tdengine batch writer"
Task 7: Production Sysctl and Service Guardrail Rollout
Files:
-
Modify:
docs/ops/vehicle-ingest-runbook.md -
Modify:
docs/ops/go-vehicle-ingest-memory.md -
Step 1: Apply sysctl changes on ECS
Run on ECS:
cat >/etc/sysctl.d/zz-lingniu-vehicle-ingest.conf <<'EOF'
net.core.somaxconn = 65535
net.ipv4.tcp_max_syn_backlog = 65535
net.ipv4.ip_local_port_range = 10000 65000
net.ipv4.tcp_tw_reuse = 1
EOF
sysctl --system
- Step 2: Verify sysctl
Run:
sysctl net.core.somaxconn net.ipv4.tcp_max_syn_backlog net.ipv4.ip_local_port_range net.ipv4.tcp_tw_reuse
Expected:
net.core.somaxconn = 65535
net.ipv4.tcp_max_syn_backlog = 65535
net.ipv4.ip_local_port_range = 10000 65000
net.ipv4.tcp_tw_reuse = 1
- Step 3: Restart gateway
Run:
systemctl restart lingniu-go-gateway.service
curl -fsS http://127.0.0.1:20211/readyz
ss -ltnp | grep -E ':(808|32960|20211) '
- Step 4: Document rollout
Add the applied sysctl values and timestamp to docs/ops/go-vehicle-ingest-memory.md.
Actual release: /opt/lingniu-go-native/releases/100k-hardening-20260703175730.
- Step 5: Commit
git add docs/ops/go-vehicle-ingest-memory.md docs/ops/vehicle-ingest-runbook.md
git commit -m "ops: document 100k gateway sysctl rollout"
Task 8: First 10K Connection Test
Files:
-
Modify:
docs/ops/100k-capacity-baseline.md -
Modify:
go/vehicle-gateway/cmd/load-sim/main.go -
Modify:
go/vehicle-gateway/internal/loadsim/* -
Step 1: Implement JT808 connection hold mode
Extended load-sim with -send=false to open N TCP connections and keep them open for duration, without sending frames.
Expected CLI:
go run ./cmd/load-sim -protocol jt808 -addr 115.29.187.205:808 -connections 10000 -send-interval 30s -duration 5m -send=false
- Step 2: Run from non-production client
Run a small smoke first:
go run ./cmd/load-sim -protocol jt808 -addr 115.29.187.205:808 -connections 100 -duration 30s -send=false
Actual: ran on ECS loopback 127.0.0.1:808 to avoid external network limits and avoid business frame pollution.
- Step 3: Capture metrics during test
On ECS:
curl -fsS http://127.0.0.1:20211/metrics | grep vehicle_gateway_active_connections
uptime
free -m
ss -tn sport = :808 | wc -l
journalctl -u lingniu-go-gateway.service --since '10 minutes ago' --no-pager | grep -Ei 'error|failed|panic|fatal' || true
- Step 4: Update capacity baseline
Record:
- max active connections reached
- CPU/load
- memory
- connection rejection count
- errors
Actual: recorded 100 / 1,000 / 10,000 / 50,000 / 100,000 total hold-only results in docs/ops/100k-capacity-baseline.md.
- Step 5: Commit
git add docs/ops/100k-capacity-baseline.md go/vehicle-gateway/cmd/load-sim go/vehicle-gateway/internal/loadsim
git commit -m "test(go): record 10k gateway connection baseline"
Self-Review
Spec coverage:
- 10W production objective: covered by capacity assumptions and phased tasks.
- Gateway connection limit and OS tuning: Tasks 3, 4, 7, 8.
- NATS/Kafka stability: baseline metrics and future batching notes; bridge-specific tuning remains a follow-up after first load test.
- Redis/MySQL projector stability: Task 5.
- TDengine write bottleneck: Task 6 sets design; implementation should follow once load evidence shows the threshold.
- ECS operations memory: Tasks 1 and 7.
Known intentional gaps:
- Full 100K test is not first. The first safe milestone is 10K connection test, then 50K, then 100K.
- Multi-ECS horizontal sharding is not implemented in this plan. It should follow after single-node evidence.
- Kafka partition changes are not included yet; changing partitions without measured FPS and consumer bottlenecks would be premature.