Files
lingniu-vehicle-ingest/go/vehicle-gateway/cmd/nats-fast-writer/main.go

454 lines
14 KiB
Go

package main
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"log/slog"
"os"
"os/signal"
"strings"
"syscall"
"time"
"github.com/nats-io/nats.go"
"github.com/redis/go-redis/v9"
_ "github.com/taosdata/driver-go/v3/taosWS"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/health"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/history"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics"
)
func main() {
logger := observability.NewLogger("nats-fast-writer")
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)
}
tdDB, err := sql.Open(cfg.TDengineDriver, cfg.TDengineDSN)
if err != nil {
logger.Error("tdengine open failed", "error", err)
os.Exit(1)
}
defer tdDB.Close()
tdDB.SetMaxOpenConns(1)
tdDB.SetMaxIdleConns(1)
if err := tdDB.PingContext(ctx); err != nil {
logger.Error("tdengine ping failed", "error", err)
os.Exit(1)
}
historyWriter := history.NewWriter(tdDB)
if cfg.TDengineEnsureSchema {
if err := historyWriter.EnsureSchema(ctx, cfg.TDengineDatabase); err != nil {
logger.Error("tdengine schema bootstrap failed", "error", err)
os.Exit(1)
}
}
if strings.TrimSpace(cfg.TDengineDatabase) != "" {
if _, err := tdDB.ExecContext(ctx, "USE "+cfg.TDengineDatabase); err != nil {
logger.Error("tdengine use database failed", "database", cfg.TDengineDatabase, "error", err)
os.Exit(1)
}
}
redisClient := redis.NewClient(&redis.Options{
Addr: cfg.RedisAddr,
Username: cfg.RedisUsername,
Password: cfg.RedisPassword,
DB: cfg.RedisDB,
})
defer redisClient.Close()
if err := redisClient.Ping(ctx).Err(); err != nil {
logger.Error("redis ping failed", "error", err)
os.Exit(1)
}
realtimeRepo := realtime.NewRepository(redisClient, realtime.Config{OnlineTTL: cfg.OnlineTTL})
registry := metrics.NewRegistry()
health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-fast-writer", []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
}},
{Name: "tdengine", Check: tdDB.PingContext},
{Name: "redis", Check: func(ctx context.Context) error { return redisClient.Ping(ctx).Err() }},
}, registry))
logger.Info("nats fast writer started",
"stream", cfg.NATSStream,
"durable", cfg.NATSDurable,
"filter", cfg.NATSFilter,
"batch_size", cfg.BatchSize,
"workers", cfg.Workers,
"operation_timeout_ms", cfg.OperationWait.Milliseconds())
for i := 0; i < cfg.Workers; i++ {
go runFastWorker(ctx, logger, registry, js, sub, historyWriter, realtimeRepo, cfg)
}
<-ctx.Done()
}
type config struct {
NATSURL string
NATSClientName string
NATSStream string
NATSDurable string
NATSFilter string
NATSSubjects []string
BatchSize int
FetchWait time.Duration
OperationWait time.Duration
AckWait time.Duration
StreamMaxAge time.Duration
Workers int
TDengineDriver string
TDengineDSN string
TDengineDatabase string
TDengineEnsureSchema bool
RedisAddr string
RedisUsername string
RedisPassword string
RedisDB int
OnlineTTL time.Duration
}
func loadConfig() config {
subjects := splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join([]string{
env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)),
env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)),
env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)),
env("NATS_SUBJECT_GB32960_FIELDS", env("KAFKA_TOPIC_GB32960_FIELDS", topics.FieldsGB32960)),
env("NATS_SUBJECT_JT808_FIELDS", env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808)),
env("NATS_SUBJECT_YUTONG_MQTT_FIELDS", env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT)),
}, ",")))
return config{
NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"),
NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-fast-writer"),
NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"),
NATSDurable: env("NATS_DURABLE", "vehicle-fast-writer"),
NATSFilter: env("NATS_FILTER", "vehicle.raw.go.>"),
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,
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second,
StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour,
Workers: envInt("FAST_WRITER_WORKERS", 8),
TDengineDriver: env("TDENGINE_DRIVER", "taosWS"),
TDengineDSN: env("TDENGINE_DSN", ""),
TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase),
TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false",
RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"),
RedisUsername: env("REDIS_USERNAME", ""),
RedisPassword: env("REDIS_PASSWORD", ""),
RedisDB: envInt("REDIS_DB", 0),
OnlineTTL: time.Duration(envInt("REALTIME_ONLINE_TTL_SECONDS", 60)) * time.Second,
}
}
type fastAppender interface {
AppendAll(context.Context, envelope.FrameEnvelope) error
AppendAllBatch(context.Context, []envelope.FrameEnvelope) error
}
type fastUpdater interface {
FastUpdate(context.Context, envelope.FrameEnvelope) error
}
type fastMessage struct {
subject string
data []byte
ack func() error
}
func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, infoReader natsConsumerInfoReader, sub natsPullSubscription, appender fastAppender, updater fastUpdater, 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 {
recordFastNATSConsumerInfoMetrics(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
}
fastMessages := make([]*fastMessage, 0, len(msgs))
for _, msg := range msgs {
natsMsg := msg
fastMessages = append(fastMessages, &fastMessage{
subject: natsMsg.Subject,
data: natsMsg.Data,
ack: func() error {
return natsMsg.Ack()
},
})
}
operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait)
err = processFastBatch(operationCtx, registry, appender, updater, fastMessages)
cancel()
if err != nil {
addFastMetric(registry, fastBatchSubject(fastMessages), "error")
logger.Error("fast write batch failed", "messages", len(fastMessages), "error", err)
continue
}
for _, msg := range fastMessages {
addFastMetric(registry, msg.subject, "ok")
}
}
}
type natsPullSubscription interface {
Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error)
}
type natsConsumerInfoReader interface {
ConsumerInfo(stream string, name string, opts ...nats.JSOpt) (*nats.ConsumerInfo, error)
}
var fastWriterStageDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000}
func processFastBatch(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, messages []*fastMessage) error {
if len(messages) == 0 {
return nil
}
envelopes := make([]envelope.FrameEnvelope, 0, len(messages))
validMessages := make([]*fastMessage, 0, len(messages))
for _, msg := range messages {
var env envelope.FrameEnvelope
if err := json.Unmarshal(msg.data, &env); err != nil {
if msg.ack != nil {
_ = msg.ack()
}
continue
}
envelopes = append(envelopes, env)
validMessages = append(validMessages, msg)
}
if len(envelopes) == 0 {
return nil
}
setFastBatchPending(registry, len(messages), len(envelopes))
defer setFastBatchPending(registry, 0, 0)
subject := fastBatchSubject(validMessages)
started := time.Now()
err := appender.AppendAllBatch(ctx, envelopes)
recordFastWriterStageDuration(registry, subject, "tdengine", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("tdengine batch append: %w", err)
}
for i, env := range envelopes {
msg := validMessages[i]
started = time.Now()
err = updater.FastUpdate(ctx, env)
recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("redis fast update: %w", err)
}
if msg.ack != nil {
started = time.Now()
err = msg.ack()
recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("nats ack: %w", err)
}
}
}
return nil
}
func processFastMessage(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage) error {
var env envelope.FrameEnvelope
if err := json.Unmarshal(msg.data, &env); err != nil {
if msg.ack != nil {
_ = msg.ack()
}
return nil
}
started := time.Now()
err := appender.AppendAll(ctx, env)
recordFastWriterStageDuration(registry, msg.subject, "tdengine", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("tdengine append: %w", err)
}
started = time.Now()
err = updater.FastUpdate(ctx, env)
recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("redis fast update: %w", err)
}
if msg.ack != nil {
started = time.Now()
err = msg.ack()
recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started))
if err != nil {
return fmt.Errorf("nats ack: %w", err)
}
}
return nil
}
func fastBatchSubject(messages []*fastMessage) string {
if len(messages) == 0 {
return "unknown"
}
subject := messages[0].subject
for _, msg := range messages[1:] {
if msg.subject != subject {
return "mixed"
}
}
if strings.TrimSpace(subject) == "" {
return "unknown"
}
return subject
}
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 addFastMetric(registry *metrics.Registry, subject string, status string) {
if registry == nil {
return
}
registry.IncCounter("vehicle_fast_writer_messages_total", metrics.Labels{"subject": subject, "status": status})
}
func setFastBatchPending(registry *metrics.Registry, messages int, envelopes int) {
if registry == nil {
return
}
registry.SetGauge("vehicle_fast_writer_batch_pending_messages", nil, float64(messages))
registry.SetGauge("vehicle_fast_writer_batch_pending_envelopes", nil, float64(envelopes))
}
func recordFastNATSConsumerInfoMetrics(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_fast_writer_nats_consumer_pending", labels, float64(info.NumPending))
registry.SetGauge("vehicle_fast_writer_nats_consumer_ack_pending", labels, float64(info.NumAckPending))
registry.SetGauge("vehicle_fast_writer_nats_consumer_waiting", labels, float64(info.NumWaiting))
}
func recordFastWriterStageDuration(registry *metrics.Registry, subject string, stage string, status string, elapsed time.Duration) {
if registry == nil {
return
}
registry.ObserveHistogram("vehicle_fast_writer_stage_duration_ms_histogram", metrics.Labels{
"subject": subject,
"stage": stage,
"status": status,
}, fastWriterStageDurationBucketsMS, float64(elapsed.Milliseconds()))
}
func statusFromError(err error) string {
if err != nil {
return "error"
}
return "ok"
}
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
}
var parsed int
if _, err := fmt.Sscanf(value, "%d", &parsed); err != nil {
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
}