Files
lingniu-vehicle-ingest/go/vehicle-gateway/cmd/stat-writer/main.go
2026-07-03 19:08:49 +08:00

194 lines
5.9 KiB
Go

package main
import (
"context"
"database/sql"
"encoding/json"
"os"
"os/signal"
"strings"
"syscall"
"time"
_ "github.com/go-sql-driver/mysql"
"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/stats"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics"
)
func main() {
logger := observability.NewLogger("vehicle-stat-writer")
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
cfg := loadConfig()
db, err := sql.Open("mysql", cfg.MySQLDSN)
if err != nil {
logger.Error("mysql open failed", "error", err)
os.Exit(1)
}
defer db.Close()
if err := db.PingContext(ctx); err != nil {
logger.Error("mysql ping failed", "error", err)
os.Exit(1)
}
registry := metrics.NewRegistry()
health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "vehicle-stat-writer", []health.Check{
{Name: "mysql", Check: db.PingContext},
}, registry))
writer := stats.NewWriter(db, cfg.Location)
if cfg.EnsureSchema {
if err := writer.EnsureSchema(ctx); err != nil {
logger.Error("mysql schema bootstrap failed", "error", err)
os.Exit(1)
}
}
reader := kafka.NewReader(kafka.ReaderConfig{
Brokers: cfg.KafkaBrokers,
GroupID: cfg.KafkaGroup,
GroupTopics: cfg.KafkaTopics,
MinBytes: 1,
MaxBytes: 10e6,
})
defer reader.Close()
logger.Info("stat writer started", "group", cfg.KafkaGroup, "topics", strings.Join(cfg.KafkaTopics, ","))
for {
message, err := reader.FetchMessage(ctx)
if err != nil {
if ctx.Err() != nil {
return
}
logger.Error("kafka fetch failed", "error", err)
continue
}
processStatMessage(ctx, logger, registry, writer, reader, message)
}
}
const kafkaMessageOperationTimeout = 30 * time.Second
type statAppender interface {
Append(context.Context, envelope.FrameEnvelope) error
}
type kafkaMessageCommitter interface {
CommitMessages(context.Context, ...kafka.Message) error
}
var statWriteDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000}
func processStatMessage(ctx context.Context, logger interface {
Error(string, ...any)
Warn(string, ...any)
}, registry *metrics.Registry, appender statAppender, committer kafkaMessageCommitter, message kafka.Message) {
messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout)
defer cancel()
addStatMetric(registry, "vehicle_stat_kafka_messages_total", message, "received")
addStatLagMetric(registry, message)
var env envelope.FrameEnvelope
if err := json.Unmarshal(message.Value, &env); err != nil {
addStatMetric(registry, "vehicle_stat_kafka_messages_total", message, "invalid_json")
logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
_ = committer.CommitMessages(messageCtx, message)
return
}
started := time.Now()
err := appender.Append(messageCtx, env)
recordStatWriteDuration(registry, message, statusFromError(err), time.Since(started))
if err != nil {
addStatMetric(registry, "vehicle_stat_writes_total", message, "error")
logger.Error("mysql append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err)
return
}
addStatMetric(registry, "vehicle_stat_writes_total", message, "ok")
if err := committer.CommitMessages(messageCtx, message); err != nil {
addStatMetric(registry, "vehicle_stat_kafka_commits_total", message, "error")
logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
return
}
addStatMetric(registry, "vehicle_stat_kafka_commits_total", message, "ok")
}
func addStatMetric(registry *metrics.Registry, name string, message kafka.Message, status string) {
if registry == nil {
return
}
registry.IncCounter(name, metrics.Labels{"topic": message.Topic, "status": status})
}
func addStatLagMetric(registry *metrics.Registry, message kafka.Message) {
if registry == nil {
return
}
registry.SetKafkaLag("vehicle_stat_kafka_lag", message.Topic, message.Partition, message.Offset, message.HighWaterMark)
}
func recordStatWriteDuration(registry *metrics.Registry, message kafka.Message, status string, elapsed time.Duration) {
if registry == nil {
return
}
registry.ObserveHistogram("vehicle_stat_write_duration_ms_histogram", metrics.Labels{
"topic": message.Topic,
"status": status,
}, statWriteDurationBucketsMS, float64(elapsed.Milliseconds()))
}
func statusFromError(err error) string {
if err != nil {
return "error"
}
return "ok"
}
type config struct {
KafkaBrokers []string
KafkaTopics []string
KafkaGroup string
MySQLDSN string
EnsureSchema bool
Location *time.Location
}
func loadConfig() config {
loc, err := time.LoadLocation(env("LOCAL_TZ", "Asia/Shanghai"))
if err != nil {
loc = time.FixedZone("Asia/Shanghai", 8*3600)
}
return config{
KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")),
KafkaTopics: splitCSV(env("KAFKA_TOPICS", strings.Join([]string{topics.FieldsGB32960, topics.FieldsJT808, topics.FieldsYutongMQTT}, ","))),
KafkaGroup: env("KAFKA_GROUP", "go-stat-writer"),
MySQLDSN: env("MYSQL_DSN", ""),
EnsureSchema: env("MYSQL_ENSURE_SCHEMA", "true") != "false",
Location: loc,
}
}
func env(key string, fallback string) string {
value := strings.TrimSpace(os.Getenv(key))
if value == "" {
return fallback
}
return value
}
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
}