143 lines
3.8 KiB
Go
143 lines
3.8 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/observability"
|
|
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats"
|
|
)
|
|
|
|
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)
|
|
}
|
|
|
|
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, 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
|
|
}
|
|
|
|
func processStatMessage(ctx context.Context, logger interface {
|
|
Error(string, ...any)
|
|
Warn(string, ...any)
|
|
}, appender statAppender, committer kafkaMessageCommitter, message kafka.Message) {
|
|
messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout)
|
|
defer cancel()
|
|
|
|
var env envelope.FrameEnvelope
|
|
if err := json.Unmarshal(message.Value, &env); err != nil {
|
|
logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
|
|
_ = committer.CommitMessages(messageCtx, message)
|
|
return
|
|
}
|
|
if err := appender.Append(messageCtx, env); err != nil {
|
|
logger.Error("mysql append failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "event_id", env.StableEventID(), "error", err)
|
|
return
|
|
}
|
|
if err := committer.CommitMessages(messageCtx, message); err != nil {
|
|
logger.Error("kafka commit failed", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err)
|
|
}
|
|
}
|
|
|
|
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", "vehicle.raw.gb32960.v1,vehicle.raw.jt808.v1")),
|
|
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
|
|
}
|