fix(gateway): guard nats stream retention and writer timeouts
This commit is contained in:
@@ -28,6 +28,7 @@ NATS 部署在 Kafka ECS 内网 `172.17.111.56:4222`。
|
||||
| Subjects | `vehicle.raw.go.gb32960.v1`、`vehicle.raw.go.jt808.v1`、`vehicle.raw.go.yutong-mqtt.v1` |
|
||||
| Durable consumer | `vehicle-kafka-bridge` |
|
||||
| 语义 | Gateway 先写 NATS raw subject;bridge 写 Kafka 成功后才 ACK NATS |
|
||||
| 保留策略 | `NATS_STREAM_MAX_BYTES=21474836480`,`NATS_STREAM_MAX_AGE_HOURS=24`,`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60`,NATS 仅做入口缓冲 |
|
||||
|
||||
### Kafka
|
||||
|
||||
|
||||
@@ -88,6 +88,19 @@ systemd gateway 已配置 `LimitNOFILE=1048576`,需要持续保持。
|
||||
| Realtime MySQL | async queue depth | 不持续增长 |
|
||||
| Stat writer | `vehicle_stat_write_duration_ms_histogram_bucket` | p99 不持续上升 |
|
||||
|
||||
## Retention Guardrails
|
||||
|
||||
10W 车辆接入时,中间件只承担解耦和短期缓冲职责,不能成为长期历史库。
|
||||
|
||||
| Component | Guardrail |
|
||||
| --- | --- |
|
||||
| NATS JetStream `VEHICLE_INGEST` | `NATS_STREAM_MAX_BYTES=21474836480`,`NATS_STREAM_MAX_AGE_HOURS=24`,`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60` |
|
||||
| Kafka `vehicle.raw.go.*` / `vehicle.fields.go.*` | `retention.ms=21600000`、`segment.ms=600000`、`segment.bytes=268435456` |
|
||||
| Fast writer | `FAST_WRITER_OPERATION_TIMEOUT_MS=1000`,避免 TDengine 批量写尾延迟导致整批重投 |
|
||||
| Root disk | 使用率低于 `80%`,超过 `85%` 进入容量告警 |
|
||||
|
||||
如果 NATS `consumer_pending` 下降但仍高,说明系统在追历史积压;如果 `ack_pending=0` 且 Kafka lag 为 `0`,不要重启服务,重点观察 NATS data size 和 pending 下降斜率。
|
||||
|
||||
## 分阶段压测
|
||||
|
||||
压测入口:
|
||||
|
||||
@@ -196,13 +196,30 @@ Kafka 是可回放日志。只要 Kafka 还保留消息,恢复 consumer 后 la
|
||||
1. 检查 bridge `/readyz` 和日志。
|
||||
2. 对比 `vehicle_bridge_kafka_writes_total` 和 `vehicle_bridge_nats_acks_total`。
|
||||
3. 如果 Kafka 写入失败,检查 ECS 到 Kafka broker 的网络。
|
||||
4. 如果 ack-pending 卡住且日志持续报错,只重启 bridge。
|
||||
4. 如果 `ack_pending > 0`,优先检查 Kafka 是否监听 `9092`、Kafka 盘是否打满,再看 bridge 日志。
|
||||
5. 如果 `ack_pending = 0` 但 `consumer_pending` 下降,说明 bridge 正在追历史积压;不要反复重启,持续观察下降速度即可。
|
||||
6. 如果 `consumer_pending` 不下降,才考虑扩容 bridge 或降低 batch/fetch 等配置。
|
||||
|
||||
```bash
|
||||
journalctl -u lingniu-go-nats-kafka-bridge.service --since '10 minutes ago' --no-pager
|
||||
systemctl restart lingniu-go-nats-kafka-bridge.service
|
||||
```
|
||||
|
||||
生产 stream 必须设置字节上限,避免 NATS JetStream 在 Kafka 或下游故障时吃满根盘:
|
||||
|
||||
```bash
|
||||
grep NATS_STREAM_MAX_BYTES /opt/lingniu-go-native/env/nats-fast-writer.env
|
||||
grep NATS_STREAM_MAX_BYTES /opt/lingniu-go-native/env/nats-kafka-bridge.env
|
||||
grep NATS_STREAM_ENSURE_TIMEOUT_SECONDS /opt/lingniu-go-native/env/nats-fast-writer.env
|
||||
grep NATS_STREAM_ENSURE_TIMEOUT_SECONDS /opt/lingniu-go-native/env/nats-kafka-bridge.env
|
||||
du -sh /opt/lingniu-nats/data
|
||||
df -h /
|
||||
```
|
||||
|
||||
当前建议值:`NATS_STREAM_MAX_BYTES=21474836480`,即 `20GiB`;`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60`,避免大 stream 元数据更新时被 NATS 客户端默认 5s 超时误杀。Kafka topic 只作为短期缓冲,当前建议保留 `6h`,不要把 Kafka 或 NATS 当长期历史存储;长期历史和 RAW 查询以 TDengine/MySQL 投影为准。
|
||||
|
||||
fast-writer 的 `FAST_WRITER_OPERATION_TIMEOUT_MS` 建议为 `1000`。实时链路仍以 100ms 级为目标,但 TDengine 批量写存在尾延迟,过小的超时会造成 NATS 消息反复重投和重复写压力。
|
||||
|
||||
### Raw 有数据但实时查不到
|
||||
|
||||
1. 查 realtime Kafka lag。
|
||||
|
||||
@@ -132,6 +132,8 @@ type config struct {
|
||||
OperationWait time.Duration
|
||||
AckWait time.Duration
|
||||
StreamMaxAge time.Duration
|
||||
StreamMaxBytes int64
|
||||
StreamEnsureWait time.Duration
|
||||
Workers int
|
||||
TDengineDriver string
|
||||
TDengineDSN string
|
||||
@@ -167,9 +169,11 @@ func loadConfig() config {
|
||||
NATSSubjects: subjects,
|
||||
BatchSize: envInt("FAST_WRITER_BATCH_SIZE", 100),
|
||||
FetchWait: time.Duration(envInt("FAST_WRITER_FETCH_WAIT_MS", 100)) * time.Millisecond,
|
||||
OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 100)) * time.Millisecond,
|
||||
OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 1000)) * time.Millisecond,
|
||||
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second,
|
||||
StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour,
|
||||
StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024),
|
||||
StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second,
|
||||
Workers: workers,
|
||||
TDengineDriver: env("TDENGINE_DRIVER", "taosWS"),
|
||||
TDengineDSN: env("TDENGINE_DSN", ""),
|
||||
@@ -381,23 +385,27 @@ func fastBatchSubject(messages []*fastMessage) string {
|
||||
}
|
||||
|
||||
func ensureStream(js nats.JetStreamContext, cfg config) error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), cfg.StreamEnsureWait)
|
||||
defer cancel()
|
||||
opts := []nats.JSOpt{nats.Context(ctx)}
|
||||
stream := &nats.StreamConfig{
|
||||
Name: cfg.NATSStream,
|
||||
Subjects: cfg.NATSSubjects,
|
||||
Storage: nats.FileStorage,
|
||||
Retention: nats.LimitsPolicy,
|
||||
MaxAge: cfg.StreamMaxAge,
|
||||
MaxBytes: cfg.StreamMaxBytes,
|
||||
Duplicates: 2 * time.Minute,
|
||||
}
|
||||
if _, err := js.StreamInfo(cfg.NATSStream); err == nil {
|
||||
_, err = js.UpdateStream(stream)
|
||||
if _, err := js.StreamInfo(cfg.NATSStream, opts...); err == nil {
|
||||
_, err = js.UpdateStream(stream, opts...)
|
||||
return err
|
||||
}
|
||||
_, err := js.AddStream(stream)
|
||||
_, err := js.AddStream(stream, opts...)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
_, updateErr := js.UpdateStream(stream)
|
||||
_, updateErr := js.UpdateStream(stream, opts...)
|
||||
if updateErr == nil {
|
||||
return nil
|
||||
}
|
||||
@@ -467,6 +475,18 @@ func envInt(key string, fallback int) int {
|
||||
return parsed
|
||||
}
|
||||
|
||||
func envInt64(key string, fallback int64) int64 {
|
||||
value := strings.TrimSpace(os.Getenv(key))
|
||||
if value == "" {
|
||||
return fallback
|
||||
}
|
||||
var parsed int64
|
||||
if _, err := fmt.Sscanf(value, "%d", &parsed); err != nil || parsed <= 0 {
|
||||
return fallback
|
||||
}
|
||||
return parsed
|
||||
}
|
||||
|
||||
func splitCSV(value string) []string {
|
||||
var out []string
|
||||
for _, item := range strings.Split(value, ",") {
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
@@ -243,6 +244,38 @@ func TestLoadConfigDefaultsTDenginePoolToSingleConnection(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigDefaultsStreamMaxBytes(t *testing.T) {
|
||||
cfg := loadConfig()
|
||||
|
||||
if got, want := cfg.StreamMaxBytes, int64(20*1024*1024*1024); got != want {
|
||||
t.Fatalf("StreamMaxBytes = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := cfg.StreamEnsureWait, 60*time.Second; got != want {
|
||||
t.Fatalf("StreamEnsureWait = %v, want %v", got, want)
|
||||
}
|
||||
if got, want := cfg.OperationWait, time.Second; got != want {
|
||||
t.Fatalf("OperationWait = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigReadsStreamMaxBytesOverride(t *testing.T) {
|
||||
t.Setenv("NATS_STREAM_MAX_BYTES", "1073741824")
|
||||
t.Setenv("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", "90")
|
||||
t.Setenv("FAST_WRITER_OPERATION_TIMEOUT_MS", "1500")
|
||||
|
||||
cfg := loadConfig()
|
||||
|
||||
if got, want := cfg.StreamMaxBytes, int64(1073741824); got != want {
|
||||
t.Fatalf("StreamMaxBytes = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := cfg.StreamEnsureWait, 90*time.Second; got != want {
|
||||
t.Fatalf("StreamEnsureWait = %v, want %v", got, want)
|
||||
}
|
||||
if got, want := cfg.OperationWait, 1500*time.Millisecond; got != want {
|
||||
t.Fatalf("OperationWait = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigReadsTDenginePoolOverride(t *testing.T) {
|
||||
t.Setenv("FAST_WRITER_WORKERS", "12")
|
||||
t.Setenv("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", "4")
|
||||
|
||||
@@ -85,19 +85,21 @@ func main() {
|
||||
}
|
||||
|
||||
type config struct {
|
||||
NATSURL string
|
||||
NATSClientName string
|
||||
NATSStream string
|
||||
NATSDurable string
|
||||
NATSFilter string
|
||||
NATSSubjects []string
|
||||
KafkaBrokers []string
|
||||
Route map[string]string
|
||||
BatchSize int
|
||||
FetchWait time.Duration
|
||||
OperationWait time.Duration
|
||||
AckWait time.Duration
|
||||
StreamMaxAge time.Duration
|
||||
NATSURL string
|
||||
NATSClientName string
|
||||
NATSStream string
|
||||
NATSDurable string
|
||||
NATSFilter string
|
||||
NATSSubjects []string
|
||||
KafkaBrokers []string
|
||||
Route map[string]string
|
||||
BatchSize int
|
||||
FetchWait time.Duration
|
||||
OperationWait time.Duration
|
||||
AckWait time.Duration
|
||||
StreamMaxAge time.Duration
|
||||
StreamMaxBytes int64
|
||||
StreamEnsureWait time.Duration
|
||||
}
|
||||
|
||||
func loadConfig() config {
|
||||
@@ -113,19 +115,21 @@ func loadConfig() config {
|
||||
route[unifiedSubject] = unifiedTopic
|
||||
}
|
||||
return config{
|
||||
NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"),
|
||||
NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-kafka-bridge"),
|
||||
NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"),
|
||||
NATSDurable: env("NATS_DURABLE", "vehicle-kafka-bridge"),
|
||||
NATSFilter: env("NATS_FILTER", "vehicle.>"),
|
||||
NATSSubjects: splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join(mapKeys(route), ","))),
|
||||
KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")),
|
||||
Route: route,
|
||||
BatchSize: envInt("BRIDGE_BATCH_SIZE", 500),
|
||||
FetchWait: time.Duration(envInt("BRIDGE_FETCH_WAIT_MS", 1000)) * time.Millisecond,
|
||||
OperationWait: time.Duration(envInt("BRIDGE_OPERATION_TIMEOUT_MS", 30000)) * time.Millisecond,
|
||||
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 60)) * time.Second,
|
||||
StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour,
|
||||
NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"),
|
||||
NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-kafka-bridge"),
|
||||
NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"),
|
||||
NATSDurable: env("NATS_DURABLE", "vehicle-kafka-bridge"),
|
||||
NATSFilter: env("NATS_FILTER", "vehicle.>"),
|
||||
NATSSubjects: splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join(mapKeys(route), ","))),
|
||||
KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")),
|
||||
Route: route,
|
||||
BatchSize: envInt("BRIDGE_BATCH_SIZE", 500),
|
||||
FetchWait: time.Duration(envInt("BRIDGE_FETCH_WAIT_MS", 1000)) * time.Millisecond,
|
||||
OperationWait: time.Duration(envInt("BRIDGE_OPERATION_TIMEOUT_MS", 30000)) * time.Millisecond,
|
||||
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 60)) * time.Second,
|
||||
StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour,
|
||||
StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024),
|
||||
StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,23 +217,27 @@ func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Regis
|
||||
}
|
||||
|
||||
func ensureStream(js nats.JetStreamContext, cfg config) error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), cfg.StreamEnsureWait)
|
||||
defer cancel()
|
||||
opts := []nats.JSOpt{nats.Context(ctx)}
|
||||
stream := &nats.StreamConfig{
|
||||
Name: cfg.NATSStream,
|
||||
Subjects: cfg.NATSSubjects,
|
||||
Storage: nats.FileStorage,
|
||||
Retention: nats.LimitsPolicy,
|
||||
MaxAge: cfg.StreamMaxAge,
|
||||
MaxBytes: cfg.StreamMaxBytes,
|
||||
Duplicates: 2 * time.Minute,
|
||||
}
|
||||
if _, err := js.StreamInfo(cfg.NATSStream); err == nil {
|
||||
_, err = js.UpdateStream(stream)
|
||||
if _, err := js.StreamInfo(cfg.NATSStream, opts...); err == nil {
|
||||
_, err = js.UpdateStream(stream, opts...)
|
||||
return err
|
||||
}
|
||||
_, err := js.AddStream(stream)
|
||||
_, err := js.AddStream(stream, opts...)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
_, updateErr := js.UpdateStream(stream)
|
||||
_, updateErr := js.UpdateStream(stream, opts...)
|
||||
if updateErr == nil {
|
||||
return nil
|
||||
}
|
||||
@@ -362,6 +370,18 @@ func envInt(key string, fallback int) int {
|
||||
return parsed
|
||||
}
|
||||
|
||||
func envInt64(key string, fallback int64) int64 {
|
||||
value := strings.TrimSpace(os.Getenv(key))
|
||||
if value == "" {
|
||||
return fallback
|
||||
}
|
||||
parsed, err := strconv.ParseInt(value, 10, 64)
|
||||
if err != nil || parsed <= 0 {
|
||||
return fallback
|
||||
}
|
||||
return parsed
|
||||
}
|
||||
|
||||
func splitCSV(value string) []string {
|
||||
var out []string
|
||||
for _, item := range strings.Split(value, ",") {
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/segmentio/kafka-go"
|
||||
@@ -177,6 +178,31 @@ func TestLoadConfigDefaultsToGoSubjectRoutes(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigDefaultsStreamMaxBytes(t *testing.T) {
|
||||
cfg := loadConfig()
|
||||
|
||||
if got, want := cfg.StreamMaxBytes, int64(20*1024*1024*1024); got != want {
|
||||
t.Fatalf("StreamMaxBytes = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := cfg.StreamEnsureWait, 60*time.Second; got != want {
|
||||
t.Fatalf("StreamEnsureWait = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigReadsStreamMaxBytesOverride(t *testing.T) {
|
||||
t.Setenv("NATS_STREAM_MAX_BYTES", "1073741824")
|
||||
t.Setenv("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", "90")
|
||||
|
||||
cfg := loadConfig()
|
||||
|
||||
if got, want := cfg.StreamMaxBytes, int64(1073741824); got != want {
|
||||
t.Fatalf("StreamMaxBytes = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := cfg.StreamEnsureWait, 90*time.Second; got != want {
|
||||
t.Fatalf("StreamEnsureWait = %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadConfigIncludesUnifiedOnlyWhenExplicitlyConfigured(t *testing.T) {
|
||||
t.Setenv("NATS_SUBJECT_UNIFIED", "vehicle.event.go.unified.v1")
|
||||
t.Setenv("KAFKA_TOPIC_UNIFIED", "vehicle.event.go.unified.v1")
|
||||
|
||||
Reference in New Issue
Block a user