diff --git a/backend/internal/config/loadconfig.go b/backend/internal/config/loadconfig.go index 817ba98..58ee8a6 100644 --- a/backend/internal/config/loadconfig.go +++ b/backend/internal/config/loadconfig.go @@ -1,17 +1,19 @@ package config import ( + "errors" "fmt" "os" - "errors" + "strconv" + "gopkg.in/yaml.v3" ) type Config struct { - Server ServerConfig `yaml:"server"` - Database DatabaseConfig `yaml:"database"` - Redis RedisConfig `yaml:"redis"` - RabbitMQ RabbitMQConfig `yaml:"rabbitmq"` + Server ServerConfig `yaml:"server"` + Database DatabaseConfig `yaml:"database"` + Redis RedisConfig `yaml:"redis"` + RabbitMQ RabbitMQConfig `yaml:"rabbitmq"` ObservabilityConfig ObservabilityConfig `yaml:"observability"` } @@ -45,10 +47,11 @@ type ObservabilityConfig struct { Pprof PprofConfig `yaml:"pprof"` } type PprofConfig struct { - Enabled bool `yaml:"enabled"` - ApiAddr string `yaml:"api_addr"` + Enabled bool `yaml:"enabled"` + ApiAddr string `yaml:"api_addr"` WorkerAddr string `yaml:"worker_addr"` } + func Load(filename string) (Config, error) { data, err := os.ReadFile(filename) if err != nil { @@ -60,9 +63,71 @@ func Load(filename string) (Config, error) { return Config{}, fmt.Errorf("parse config %s: %w", filename, err) } + ApplyEnvOverrides(&cfg) return cfg, nil } +func ApplyEnvOverrides(cfg *Config) { + if cfg == nil { + return + } + if v := os.Getenv("SERVER_PORT"); v != "" { + if port, err := strconv.Atoi(v); err == nil { + cfg.Server.Port = port + } + } + if v := os.Getenv("MYSQL_HOST"); v != "" { + cfg.Database.Host = v + } + if v := os.Getenv("MYSQL_PORT"); v != "" { + if port, err := strconv.Atoi(v); err == nil { + cfg.Database.Port = port + } + } + if v := os.Getenv("MYSQL_USER"); v != "" { + cfg.Database.User = v + } + if v := os.Getenv("MYSQL_ROOT_PASSWORD"); v != "" { + cfg.Database.Password = v + } + if v := os.Getenv("MYSQL_PASSWORD"); v != "" { + cfg.Database.Password = v + } + if v := os.Getenv("MYSQL_DATABASE"); v != "" { + cfg.Database.DBName = v + } + if v := os.Getenv("REDIS_HOST"); v != "" { + cfg.Redis.Host = v + } + if v := os.Getenv("REDIS_PORT"); v != "" { + if port, err := strconv.Atoi(v); err == nil { + cfg.Redis.Port = port + } + } + if v := os.Getenv("REDIS_PASSWORD"); v != "" { + cfg.Redis.Password = v + } + if v := os.Getenv("REDIS_DB"); v != "" { + if db, err := strconv.Atoi(v); err == nil { + cfg.Redis.DB = db + } + } + if v := os.Getenv("RABBITMQ_HOST"); v != "" { + cfg.RabbitMQ.Host = v + } + if v := os.Getenv("RABBITMQ_PORT"); v != "" { + if port, err := strconv.Atoi(v); err == nil { + cfg.RabbitMQ.Port = port + } + } + if v := os.Getenv("RABBITMQ_USER"); v != "" { + cfg.RabbitMQ.Username = v + } + if v := os.Getenv("RABBITMQ_PASS"); v != "" { + cfg.RabbitMQ.Password = v + } +} + // bool用来表示是否使用了默认配置,true表示使用了默认配置 func LoadLocalDev(filename string) (Config, bool, error) { cfg, err := Load(filename) @@ -76,14 +141,14 @@ func LoadLocalDev(filename string) (Config, bool, error) { } func DefaultLocalConfig() Config { - return Config{ + cfg := Config{ Server: ServerConfig{ Port: 8080, }, Database: DatabaseConfig{ Host: "localhost", Port: 3306, - User: "root", + User: "root", Password: "123456", DBName: "feedsystem", }, @@ -101,10 +166,12 @@ func DefaultLocalConfig() Config { }, ObservabilityConfig: ObservabilityConfig{ Pprof: PprofConfig{ - Enabled: true, - ApiAddr: "localhost:6060", + Enabled: true, + ApiAddr: "localhost:6060", WorkerAddr: "localhost:6061", }, }, } -} \ No newline at end of file + ApplyEnvOverrides(&cfg) + return cfg +} diff --git a/backend/internal/worker/outboxworker.go b/backend/internal/worker/outboxworker.go index 1a45b20..6bc27eb 100644 --- a/backend/internal/worker/outboxworker.go +++ b/backend/internal/worker/outboxworker.go @@ -1,93 +1,107 @@ -package worker - -import ( - "context" - "encoding/json" - "feedsystem_video_go/internal/middleware/rabbitmq" - "feedsystem_video_go/internal/middleware/redis" - "feedsystem_video_go/internal/video" - "fmt" - "log" - "time" - - oredis "github.com/redis/go-redis/v9" - "gorm.io/gorm" -) - -func StartOutboxPoller(db *gorm.DB, tmq *rabbitmq.TimelineMQ) { - go func() { - for { - var messages []video.OutboxMsg - - err := db.Where("status = ?", "pending").Order("create_time ASC").Limit(100).Find(&messages).Error - - if err != nil || len(messages) == 0 { - time.Sleep(1 * time.Second) - continue - } - - for _, msg := range messages { - err := tmq.PublishVideo(context.Background(), msg.VideoID, msg.CreateTime) - - if err == nil { - db.Delete(&msg) - } else { - log.Printf("投递MQ失败: VideoID: %d, err: %v", msg.VideoID, err) - } - } - } - }() -} - -func StartConsumer(tmq *rabbitmq.TimelineMQ, queueName string, redisClient *redis.Client) { - msgs, err := tmq.Ch.Consume( - queueName, - "", - false, - false, - false, - false, - nil, - ) - - if err != nil { - log.Printf("注册消费失败") - return - } - - go func() { - for msg := range msgs { - var event rabbitmq.TimelineEvent - err := json.Unmarshal(msg.Body, &event) - - if err != nil { - log.Printf("反序列化失败") - msg.Ack(false) - continue - } - - ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) - timelineKey := redisClient.Key("feed:global_timeline") - err = redisClient.ZAdd(ctx, timelineKey, oredis.Z{ - Score: float64(event.CreateTime), - Member: fmt.Sprintf("%d", event.VideoID), - }) - - if err != nil { - log.Printf("写入Zset失败") - msg.Nack(false, true) - cancel() - continue - } - - err = redisClient.ZRemRangeByRank(ctx, timelineKey, 0, -1001) - - if err != nil { - log.Printf("ZRem失败") - } - - msg.Ack(false) - cancel() - } - }() -} +package worker + +import ( + "context" + "encoding/json" + "feedsystem_video_go/internal/middleware/rabbitmq" + "feedsystem_video_go/internal/middleware/redis" + "feedsystem_video_go/internal/video" + "fmt" + "log" + "time" + + oredis "github.com/redis/go-redis/v9" + "gorm.io/gorm" +) + +func StartOutboxPoller(db *gorm.DB, tmq *rabbitmq.TimelineMQ) { + if db == nil || tmq == nil || tmq.RabbitMQ == nil || tmq.Ch == nil { + log.Printf("Outbox poller disabled: timeline mq is not initialized") + return + } + + go func() { + for { + var messages []video.OutboxMsg + + err := db.Where("status = ?", "pending").Order("create_time ASC").Limit(100).Find(&messages).Error + + if err != nil || len(messages) == 0 { + time.Sleep(1 * time.Second) + continue + } + + for _, msg := range messages { + err := tmq.PublishVideo(context.Background(), msg.VideoID, msg.CreateTime) + + if err == nil { + db.Delete(&msg) + } else { + log.Printf("投递MQ失败: VideoID: %d, err: %v", msg.VideoID, err) + } + } + } + }() +} + +func StartConsumer(tmq *rabbitmq.TimelineMQ, queueName string, redisClient *redis.Client) { + if tmq == nil || tmq.RabbitMQ == nil || tmq.Ch == nil { + log.Printf("Timeline consumer disabled: timeline mq is not initialized") + return + } + if redisClient == nil { + log.Printf("Timeline consumer disabled: redis is not initialized") + return + } + + msgs, err := tmq.Ch.Consume( + queueName, + "", + false, + false, + false, + false, + nil, + ) + + if err != nil { + log.Printf("注册消费失败") + return + } + + go func() { + for msg := range msgs { + var event rabbitmq.TimelineEvent + err := json.Unmarshal(msg.Body, &event) + + if err != nil { + log.Printf("反序列化失败") + msg.Ack(false) + continue + } + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + timelineKey := redisClient.Key("feed:global_timeline") + err = redisClient.ZAdd(ctx, timelineKey, oredis.Z{ + Score: float64(event.CreateTime), + Member: fmt.Sprintf("%d", event.VideoID), + }) + + if err != nil { + log.Printf("写入Zset失败") + msg.Nack(false, true) + cancel() + continue + } + + err = redisClient.ZRemRangeByRank(ctx, timelineKey, 0, -1001) + + if err != nil { + log.Printf("ZRem失败") + } + + msg.Ack(false) + cancel() + } + }() +} diff --git a/docker-compose.yml b/docker-compose.yml index 358660a..413a8b8 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -16,22 +16,24 @@ services: - --default-authentication-plugin=mysql_native_password - --character-set-server=utf8mb4 - --collation-server=utf8mb4_0900_ai_ci - healthcheck: - test: ["CMD-SHELL", "mysqladmin ping -h 127.0.0.1 -uroot -p123456 --silent"] + healthcheck: + test: ["CMD-SHELL", "mysqladmin ping -h 127.0.0.1 -uroot -p$${MYSQL_ROOT_PASSWORD} --silent"] interval: 5s timeout: 5s retries: 20 - redis: - image: redis:7-alpine - restart: always - command: ["redis-server", "--appendonly", "yes", "--requirepass", "${REDIS_PASSWORD:-123456}"] + redis: + image: redis:7-alpine + restart: always + environment: + REDIS_PASSWORD: ${REDIS_PASSWORD:-123456} + command: ["redis-server", "--appendonly", "yes", "--requirepass", "${REDIS_PASSWORD:-123456}"] ports: - "6379:6379" volumes: - redis_data:/data - healthcheck: - test: ["CMD", "redis-cli", "-a", "123456", "ping"] + healthcheck: + test: ["CMD-SHELL", "redis-cli -a \"$${REDIS_PASSWORD}\" ping"] interval: 5s timeout: 3s retries: 20 @@ -59,8 +61,13 @@ services: dockerfile: backend/Dockerfile target: api restart: always - environment: - JWT_SECRET: ${JWT_SECRET:-feedsystem-dev-secret-key} + environment: + JWT_SECRET: ${JWT_SECRET:-feedsystem-dev-secret-key} + MYSQL_DATABASE: ${MYSQL_DATABASE:-feedsystem} + MYSQL_ROOT_PASSWORD: ${MYSQL_ROOT_PASSWORD:-123456} + REDIS_PASSWORD: ${REDIS_PASSWORD:-123456} + RABBITMQ_USER: ${RABBITMQ_USER:-admin} + RABBITMQ_PASS: ${RABBITMQ_PASS:-password123} ports: - "8080:8080" volumes: @@ -85,8 +92,13 @@ services: dockerfile: backend/Dockerfile target: worker restart: always - environment: - JWT_SECRET: ${JWT_SECRET:-feedsystem-dev-secret-key} + environment: + JWT_SECRET: ${JWT_SECRET:-feedsystem-dev-secret-key} + MYSQL_DATABASE: ${MYSQL_DATABASE:-feedsystem} + MYSQL_ROOT_PASSWORD: ${MYSQL_ROOT_PASSWORD:-123456} + REDIS_PASSWORD: ${REDIS_PASSWORD:-123456} + RABBITMQ_USER: ${RABBITMQ_USER:-admin} + RABBITMQ_PASS: ${RABBITMQ_PASS:-password123} volumes: - ./backend/configs/config.docker.yaml:/app/configs/config.yaml:ro depends_on: