refactor: replace Kafka topic configuration with mqbuild constants for improved readability

pull/3759/head
withchao 2 months ago
parent fa411a79a4
commit bfe8bcd9ad

@ -498,6 +498,7 @@ func startRedisServerRegister(ctx context.Context, cfg *serverConfig, client dis
log.ZWarn(ctx, "gateway unregister failed", err, "address", selfAddr) log.ZWarn(ctx, "gateway unregister failed", err, "address", selfAddr)
} }
}() }()
register()
for { for {
select { select {
case <-timer.C: case <-timer.C:

@ -73,11 +73,11 @@ func Start(ctx context.Context, config *Config, client discovery.SvcDiscoveryReg
if err != nil { if err != nil {
return err return err
} }
mongoProducer, err := builder.GetTopicProducer(ctx, config.KafkaConfig.ToMongoTopic) mongoProducer, err := builder.GetTopicProducer(ctx, mqbuild.TopicToMongo)
if err != nil { if err != nil {
return err return err
} }
pushProducer, err := builder.GetTopicProducer(ctx, config.KafkaConfig.ToPushTopic) pushProducer, err := builder.GetTopicProducer(ctx, mqbuild.TopicToPush)
if err != nil { if err != nil {
return err return err
} }
@ -100,11 +100,11 @@ func Start(ctx context.Context, config *Config, client discovery.SvcDiscoveryReg
if err != nil { if err != nil {
return err return err
} }
historyConsumer, err := builder.GetTopicConsumer(ctx, config.KafkaConfig.ToRedisTopic) historyConsumer, err := builder.GetTopicConsumer(ctx, mqbuild.TopicToRedis)
if err != nil { if err != nil {
return err return err
} }
historyMongoConsumer, err := builder.GetTopicConsumer(ctx, config.KafkaConfig.ToMongoTopic) historyMongoConsumer, err := builder.GetTopicConsumer(ctx, mqbuild.TopicToMongo)
if err != nil { if err != nil {
return err return err
} }

@ -65,17 +65,17 @@ func Start(ctx context.Context, config *Config, client discovery.SvcDiscoveryReg
if err != nil { if err != nil {
return err return err
} }
offlinePushProducer, err := builder.GetTopicProducer(ctx, config.KafkaConfig.ToOfflinePushTopic) offlinePushProducer, err := builder.GetTopicProducer(ctx, mqbuild.TopicToOfflinePush)
if err != nil { if err != nil {
return err return err
} }
database := controller.NewPushDatabase(cacheModel, offlinePushProducer) database := controller.NewPushDatabase(cacheModel, offlinePushProducer)
pushConsumer, err := builder.GetTopicConsumer(ctx, config.KafkaConfig.ToPushTopic) pushConsumer, err := builder.GetTopicConsumer(ctx, mqbuild.TopicToPush)
if err != nil { if err != nil {
return err return err
} }
offlinePushConsumer, err := builder.GetTopicConsumer(ctx, config.KafkaConfig.ToOfflinePushTopic) offlinePushConsumer, err := builder.GetTopicConsumer(ctx, mqbuild.TopicToOfflinePush)
if err != nil { if err != nil {
return err return err
} }

@ -92,7 +92,7 @@ func Start(ctx context.Context, config *Config, client discovery.SvcDiscoveryReg
if err != nil { if err != nil {
return err return err
} }
redisProducer, err := builder.GetTopicProducer(ctx, config.KafkaConfig.ToRedisTopic) redisProducer, err := builder.GetTopicProducer(ctx, mqbuild.TopicToRedis)
if err != nil { if err != nil {
return err return err
} }

@ -14,7 +14,7 @@ const (
func NormalizeQueueEngine(engine string) string { func NormalizeQueueEngine(engine string) string {
switch strings.ToLower(strings.TrimSpace(engine)) { switch strings.ToLower(strings.TrimSpace(engine)) {
case "kafka": case "", "kafka":
return QueueEngineKafka return QueueEngineKafka
case "redis": case "redis":
return QueueEngineRedis return QueueEngineRedis

Loading…
Cancel
Save