package main import ( "context" "encoding/json" "errors" "fmt" "log/slog" "os" "os/signal" "strconv" "strings" "syscall" "time" "github.com/nats-io/nats.go" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" ) func main() { logger := observability.NewLogger("nats-kafka-bridge") ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() cfg := loadConfig() conn, err := nats.Connect(cfg.NATSURL, nats.Name(cfg.NATSClientName), nats.Timeout(5*time.Second)) if err != nil { logger.Error("nats connect failed", "error", err) os.Exit(1) } defer conn.Close() registry := metrics.NewRegistry() health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-kafka-bridge", []health.Check{ {Name: "nats", Check: func(context.Context) error { if conn.Status() != nats.CONNECTED { return fmt.Errorf("nats status is %s", conn.Status().String()) } return nil }}, }, registry)) js, err := conn.JetStream() if err != nil { logger.Error("nats jetstream init failed", "error", err) os.Exit(1) } if err := ensureStream(js, cfg); err != nil { logger.Error("nats stream ensure failed", "stream", cfg.NATSStream, "error", err) os.Exit(1) } sub, err := js.PullSubscribe( cfg.NATSFilter, cfg.NATSDurable, nats.BindStream(cfg.NATSStream), nats.ManualAck(), nats.AckWait(cfg.AckWait), nats.MaxDeliver(-1), ) if err != nil { logger.Error("nats pull consumer init failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) os.Exit(1) } writer := &kafka.Writer{ Addr: kafka.TCP(cfg.KafkaBrokers...), Balancer: &kafka.Hash{}, AllowAutoTopicCreation: false, RequiredAcks: kafka.RequireAll, Async: false, } defer writer.Close() logger.Info("nats kafka bridge started", "nats_url", cfg.NATSURL, "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "filter", cfg.NATSFilter, "kafka_brokers", strings.Join(cfg.KafkaBrokers, ",")) runBridge(ctx, logger, registry, js, sub, writer, cfg) } 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 StreamMaxBytes int64 StreamEnsureWait time.Duration } func loadConfig() config { route := map[string]string{ env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)): env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960), env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)): env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808), env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)): env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT), env("NATS_SUBJECT_GB32960_FIELDS", env("KAFKA_TOPIC_GB32960_FIELDS", topics.FieldsGB32960)): env("KAFKA_TOPIC_GB32960_FIELDS", topics.FieldsGB32960), env("NATS_SUBJECT_JT808_FIELDS", env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808)): env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808), env("NATS_SUBJECT_YUTONG_MQTT_FIELDS", env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT)): env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT), } if unifiedSubject, unifiedTopic, ok := unifiedRouteFromEnv(); ok { 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, StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024), StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second, } } func unifiedRouteFromEnv() (string, string, bool) { subject, hasSubject := envOptional("NATS_SUBJECT_UNIFIED") topic, hasTopic := envOptional("KAFKA_TOPIC_UNIFIED") if !hasSubject && !hasTopic { return "", "", false } if subject == "" { subject = topic } if topic == "" { topic = subject } if subject == "" || topic == "" { return "", "", false } return subject, topic, true } type natsPullSubscription interface { Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error) } type natsConsumerInfoReader interface { ConsumerInfo(stream string, name string, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) } type kafkaBatchWriter interface { WriteMessages(context.Context, ...kafka.Message) error } type bridgeMessage struct { subject string data []byte ack func() error } func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, infoReader natsConsumerInfoReader, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { var lastConsumerInfoAt time.Time for { if ctx.Err() != nil { return } if time.Since(lastConsumerInfoAt) >= time.Duration(envInt("NATS_CONSUMER_METRICS_INTERVAL_SECONDS", 10))*time.Second { lastConsumerInfoAt = time.Now() if infoReader != nil { info, err := infoReader.ConsumerInfo(cfg.NATSStream, cfg.NATSDurable, nats.Context(ctx)) if err != nil { logger.Warn("nats consumer info failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) } else { recordNATSConsumerInfoMetrics(registry, cfg, info) } } } msgs, err := sub.Fetch(cfg.BatchSize, nats.MaxWait(cfg.FetchWait)) if err != nil { if errors.Is(err, nats.ErrTimeout) { continue } logger.Error("nats fetch failed", "error", err) time.Sleep(time.Second) continue } bridgeMessages := make([]bridgeMessage, 0, len(msgs)) for _, msg := range msgs { msg := msg bridgeMessages = append(bridgeMessages, bridgeMessage{ subject: msg.Subject, data: msg.Data, ack: func() error { return msg.Ack() }, }) } operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) err = bridgeBatch(operationCtx, registry, writer, bridgeMessages, cfg.Route) cancel() if err != nil { logger.Error("bridge batch failed", "count", len(bridgeMessages), "error", err) continue } } } 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, opts...); err == nil { _, err = js.UpdateStream(stream, opts...) return err } _, err := js.AddStream(stream, opts...) if err == nil { return nil } _, updateErr := js.UpdateStream(stream, opts...) if updateErr == nil { return nil } return err } func bridgeBatch(ctx context.Context, registry *metrics.Registry, writer kafkaBatchWriter, messages []bridgeMessage, route map[string]string) error { if len(messages) == 0 { return nil } started := time.Now() status := "ok" defer func() { recordBridgeBatchDuration(registry, status, time.Since(started)) recordBridgeBatchPending(registry, 0, 0) }() kafkaMessages := make([]kafka.Message, 0, len(messages)) for _, message := range messages { addBridgeSubjectMetric(registry, "vehicle_bridge_messages_total", message.subject, "received") topic, ok := route[message.subject] if !ok || topic == "" { addBridgeSubjectMetric(registry, "vehicle_bridge_messages_total", message.subject, "route_error") status = "error" return fmt.Errorf("kafka topic not configured for nats subject %q", message.subject) } kafkaMessages = append(kafkaMessages, kafkaMessage(topic, message.data)) } recordBridgeBatchPending(registry, len(messages), len(kafkaMessages)) if err := writer.WriteMessages(ctx, kafkaMessages...); err != nil { for _, message := range kafkaMessages { addBridgeTopicMetric(registry, "vehicle_bridge_kafka_writes_total", message.Topic, "error") } status = "error" return err } for _, message := range kafkaMessages { addBridgeTopicMetric(registry, "vehicle_bridge_kafka_writes_total", message.Topic, "ok") } for _, message := range messages { if message.ack == nil { continue } if err := message.ack(); err != nil { addBridgeSubjectMetric(registry, "vehicle_bridge_nats_acks_total", message.subject, "error") status = "error" return err } addBridgeSubjectMetric(registry, "vehicle_bridge_nats_acks_total", message.subject, "ok") } return nil } var bridgeBatchDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000} func addBridgeSubjectMetric(registry *metrics.Registry, name string, subject string, status string) { if registry == nil { return } registry.IncCounter(name, metrics.Labels{"subject": subject, "status": status}) } func addBridgeTopicMetric(registry *metrics.Registry, name string, topic string, status string) { if registry == nil { return } registry.IncCounter(name, metrics.Labels{"topic": topic, "status": status}) } func recordBridgeBatchPending(registry *metrics.Registry, messages int, kafkaMessages int) { if registry == nil { return } registry.SetGauge("vehicle_bridge_batch_pending_messages", nil, float64(messages)) registry.SetGauge("vehicle_bridge_batch_pending_kafka_messages", nil, float64(kafkaMessages)) } func recordBridgeBatchDuration(registry *metrics.Registry, status string, elapsed time.Duration) { if registry == nil { return } registry.ObserveHistogram("vehicle_bridge_batch_duration_ms_histogram", metrics.Labels{ "status": status, }, bridgeBatchDurationBucketsMS, float64(elapsed.Milliseconds())) } func recordNATSConsumerInfoMetrics(registry *metrics.Registry, cfg config, info *nats.ConsumerInfo) { if registry == nil || info == nil { return } labels := metrics.Labels{"stream": cfg.NATSStream, "consumer": cfg.NATSDurable} registry.SetGauge("vehicle_bridge_nats_consumer_pending", labels, float64(info.NumPending)) registry.SetGauge("vehicle_bridge_nats_consumer_ack_pending", labels, float64(info.NumAckPending)) registry.SetGauge("vehicle_bridge_nats_consumer_waiting", labels, float64(info.NumWaiting)) } func kafkaMessage(topic string, data []byte) kafka.Message { var env envelope.FrameEnvelope message := kafka.Message{Topic: topic, Value: data} if err := json.Unmarshal(data, &env); err == nil { message.Key = env.KafkaKey() } return message } func env(key string, fallback string) string { value := strings.TrimSpace(os.Getenv(key)) if value == "" { return fallback } return value } func envOptional(key string) (string, bool) { value, ok := os.LookupEnv(key) if !ok { return "", false } return strings.TrimSpace(value), true } func envInt(key string, fallback int) int { value := strings.TrimSpace(os.Getenv(key)) if value == "" { return fallback } parsed, err := strconv.Atoi(value) if err != nil || parsed <= 0 { return fallback } 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, ",") { item = strings.TrimSpace(item) if item != "" { out = append(out, item) } } return out } func mapKeys(values map[string]string) []string { out := make([]string, 0, len(values)) for key := range values { out = append(out, key) } return out }