正在显示
4 个修改的文件
包含
9 行增加
和
20 行删除
| 1 | package configs | 1 | package configs |
| 2 | 2 | ||
| 3 | import ( | 3 | import ( |
| 4 | - "os" | ||
| 5 | - "strings" | ||
| 6 | - | ||
| 7 | "gitlab.fjmaimaimai.com/mmm-go/partnermg/pkg/constant" | 4 | "gitlab.fjmaimaimai.com/mmm-go/partnermg/pkg/constant" |
| 8 | ) | 5 | ) |
| 9 | 6 | ||
| @@ -17,17 +14,5 @@ var Cfg = MqConfig{ | @@ -17,17 +14,5 @@ var Cfg = MqConfig{ | ||
| 17 | ConsumerId: constant.KafkaCfg.ConsumerId, | 14 | ConsumerId: constant.KafkaCfg.ConsumerId, |
| 18 | } | 15 | } |
| 19 | 16 | ||
| 20 | -func init() { | ||
| 21 | - | ||
| 22 | - Cfg = MqConfig{ | ||
| 23 | - Servers: []string{"106.52.15.41:9092"}, | ||
| 24 | - ConsumerId: "partnermg", | ||
| 25 | - } | ||
| 26 | - if os.Getenv("KAFKA_HOST") != "" { | ||
| 27 | - kafkaHost := os.Getenv("KAFKA_HOST") | ||
| 28 | - Cfg.Servers = strings.Split(kafkaHost, ";") | ||
| 29 | - } | ||
| 30 | -} | ||
| 31 | - | ||
| 32 | // "192.168.190.136:9092", | 17 | // "192.168.190.136:9092", |
| 33 | // "106.52.15.41:9092" | 18 | // "106.52.15.41:9092" |
| @@ -59,9 +59,9 @@ func (c *MessageConsumer) ConsumeClaim(groupSession sarama.ConsumerGroupSession, | @@ -59,9 +59,9 @@ func (c *MessageConsumer) ConsumeClaim(groupSession sarama.ConsumerGroupSession, | ||
| 59 | } | 59 | } |
| 60 | if err = topicHandle(message); err != nil { | 60 | if err = topicHandle(message); err != nil { |
| 61 | logs.Error("Message claimed: kafka消息处理错误 topic =", message.Topic, message.Offset, err) | 61 | logs.Error("Message claimed: kafka消息处理错误 topic =", message.Topic, message.Offset, err) |
| 62 | - } else { | ||
| 63 | - groupSession.MarkMessage(message, "") | ||
| 64 | } | 62 | } |
| 63 | + groupSession.MarkMessage(message, "") | ||
| 64 | + | ||
| 65 | } | 65 | } |
| 66 | return nil | 66 | return nil |
| 67 | } | 67 | } |
| @@ -20,13 +20,14 @@ func SyncBestshopOrder(message *sarama.ConsumerMessage) error { | @@ -20,13 +20,14 @@ func SyncBestshopOrder(message *sarama.ConsumerMessage) error { | ||
| 20 | ) | 20 | ) |
| 21 | err = json.Unmarshal(message.Value, &cmd) | 21 | err = json.Unmarshal(message.Value, &cmd) |
| 22 | if err != nil { | 22 | if err != nil { |
| 23 | - return fmt.Errorf("[SyncBestshopOrder] 解析kafka数据失败;%s", err) | 23 | + return fmt.Errorf("[Consumer][SyncBestshopOrder] 解析kafka数据失败;%s", err) |
| 24 | } | 24 | } |
| 25 | if cmd.PartnerId <= 0 { | 25 | if cmd.PartnerId <= 0 { |
| 26 | - logs.Info("[SyncBestshopOrder] PartnerId<=0 ,不处理消息") | 26 | + logs.Info("[Consumer][SyncBestshopOrder] PartnerId<=0 ,不处理消息") |
| 27 | return nil | 27 | return nil |
| 28 | } | 28 | } |
| 29 | srv := syncOrderSrv.NewOrderInfoService(nil) | 29 | srv := syncOrderSrv.NewOrderInfoService(nil) |
| 30 | err = srv.SyncOrderFromBestshop(cmd) | 30 | err = srv.SyncOrderFromBestshop(cmd) |
| 31 | - return err | 31 | + e := fmt.Errorf("[Consumer][SyncBestshopOrder] %s", err) |
| 32 | + return e | ||
| 32 | } | 33 | } |
-
请 注册 或 登录 后发表评论