feat(go): add nats kafka ingest bridge

This commit is contained in:
lingniu
2026-07-02 15:07:44 +08:00
parent 079bd76d38
commit c27150cc15
11 changed files with 1121 additions and 12 deletions

View File

@@ -145,6 +145,33 @@ func buildIdentityResolver(ctx context.Context, logger *slog.Logger) (identity.R
}
func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) {
if strings.TrimSpace(os.Getenv("NATS_URL")) != "" {
sink, err := eventbus.NewNATSSink(natsSinkConfigFromEnv())
if err != nil {
return nil, err
}
var out eventbus.Sink = eventbus.NewRetryingSink(sink, eventbus.RetryConfig{
Attempts: envInt("NATS_PUBLISH_ATTEMPTS", 3),
Backoff: time.Duration(envInt("NATS_PUBLISH_BACKOFF_MS", 100)) * time.Millisecond,
AttemptTimeout: time.Duration(envInt("NATS_PUBLISH_TIMEOUT_MS", 3000)) * time.Millisecond,
})
if envBool("NATS_ASYNC_ENABLED", true) {
queueSize := envInt("NATS_ASYNC_QUEUE_SIZE", envInt("KAFKA_ASYNC_QUEUE_SIZE", 100000))
workers := envInt("NATS_ASYNC_WORKERS", envInt("KAFKA_ASYNC_WORKERS", 8))
timeout := time.Duration(envInt("NATS_ASYNC_PUBLISH_TIMEOUT_MS", envInt("KAFKA_ASYNC_PUBLISH_TIMEOUT_MS", 30000))) * time.Millisecond
out = eventbus.NewAsyncSink(out, eventbus.AsyncConfig{
QueueSize: queueSize,
Workers: workers,
OperationTimeout: timeout,
OnError: func(err error) {
logger.Warn("nats async publish failed", "error", err)
},
})
logger.Info("nats async publish enabled", "queue_size", queueSize, "workers", workers, "publish_timeout_ms", timeout.Milliseconds())
}
logger.Info("nats jetstream sink enabled", "url", env("NATS_URL", ""), "unified_subject", natsSinkConfigFromEnv().UnifiedSubject)
return out, nil
}
brokers := splitCSV(os.Getenv("KAFKA_BROKERS"))
if len(brokers) == 0 {
logger.Warn("KAFKA_BROKERS is empty; using log sink")
@@ -195,6 +222,19 @@ func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error)
return out, nil
}
func natsSinkConfigFromEnv() eventbus.NATSConfig {
return eventbus.NATSConfig{
URL: env("NATS_URL", ""),
Name: env("NATS_CLIENT_NAME", "lingniu-vehicle-gateway"),
RawSubjects: map[envelope.Protocol]string{
envelope.ProtocolGB32960: env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1")),
envelope.ProtocolJT808: env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1")),
envelope.ProtocolYutongMQTT: env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1")),
},
UnifiedSubject: env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1")),
}
}
func env(key string, fallback string) string {
value := strings.TrimSpace(os.Getenv(key))
if value == "" {