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/observability" ) 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() 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, 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 } func loadConfig() config { route := map[string]string{ env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1")): env("KAFKA_TOPIC_GB32960_RAW", "vehicle.raw.gb32960.v1"), env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1")): env("KAFKA_TOPIC_JT808_RAW", "vehicle.raw.jt808.v1"), env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1")): env("KAFKA_TOPIC_YUTONG_MQTT_RAW", "vehicle.raw.yutong-mqtt.v1"), env("NATS_SUBJECT_UNIFIED", env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1")): env("KAFKA_TOPIC_UNIFIED", "vehicle.event.unified.v1"), } 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, } } type natsPullSubscription interface { Fetch(int, ...nats.PullOpt) ([]*nats.Msg, 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, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { for { if ctx.Err() != nil { return } 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, 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 { stream := &nats.StreamConfig{ Name: cfg.NATSStream, Subjects: cfg.NATSSubjects, Storage: nats.FileStorage, Retention: nats.LimitsPolicy, MaxAge: cfg.StreamMaxAge, Duplicates: 2 * time.Minute, } if _, err := js.StreamInfo(cfg.NATSStream); err == nil { _, err = js.UpdateStream(stream) return err } _, err := js.AddStream(stream) if err == nil { return nil } _, updateErr := js.UpdateStream(stream) if updateErr == nil { return nil } return err } func bridgeBatch(ctx context.Context, writer kafkaBatchWriter, messages []bridgeMessage, route map[string]string) error { if len(messages) == 0 { return nil } kafkaMessages := make([]kafka.Message, 0, len(messages)) for _, message := range messages { topic, ok := route[message.subject] if !ok || topic == "" { return fmt.Errorf("kafka topic not configured for nats subject %q", message.subject) } kafkaMessages = append(kafkaMessages, kafkaMessage(topic, message.data)) } if err := writer.WriteMessages(ctx, kafkaMessages...); err != nil { return err } for _, message := range messages { if message.ack == nil { continue } if err := message.ack(); err != nil { return err } } return nil } 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 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 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 }