# consumer_manage_pkg **Repository Path**: sizzflair/consumer_manage_pkg ## Basic Information - **Project Name**: consumer_manage_pkg - **Description**: No description available - **Primary Language**: Go - **License**: Apache-2.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2025-08-05 - **Last Updated**: 2025-12-29 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # NSQ Consumer Dynamic Management Module ## 功能特性 - **动态调整**: 支持运行时动态调整NSQ消费者数量 - **状态监控**: 实时监控消费者状态并上报到数据库 - **HTTP接口**: 提供RESTful API进行管理操作 - **配置灵活**: 支持环境变量和代码配置 - **接口驱动**: 基于接口设计,易于扩展和测试 - **生产就绪**: 包含完整的错误处理和日志记录 ## 快速开始 ### 基础集成 ```go package main import ( "database/sql" "log" "time" "your-project/app/pkg/consumer_manager" "your-project/internal/pkg/nsq" ) func main() { // 1. 准备依赖 db, _ := sql.Open("mysql", "your-dsn") nsqStruct := &nsq.NsqStruct{} // 你的NSQ结构体 // 2. 创建配置 config := consumer_manager.NewConfig() // 3. 创建管理器 manager := consumer_manager.NewManager(config, nsqStruct, db, nil) // 4. 设置回调函数 manager.SetCallFunc(yourCallbackFunction) // 5. 启动管理器 if err := manager.Start(); err != nil { log.Fatal(err) } // 6. 启动HTTP服务器(可选) apiServer := consumer_manager.NewAPIServer(manager, config, nil) go apiServer.StartServer() // 保持运行 select {} } ``` ### 集成到现有Gin服务 ```go func setupRoutes(engine *gin.Engine, manager *consumer_manager.Manager, config *consumer_manager.Config) { apiServer := consumer_manager.NewAPIServer(manager, config, nil) apiServer.SetupRoutes(engine) } ``` ## 集成到weixin_open_event_consumer的详细步骤 ### 步骤1: 修改init.go文件 在 `app/internal/app/weixin_open_event_consumer/init.go` 中添加以下代码: ```go package weixin_open_event_consumer import ( "flag" "os" "path/filepath" "zk_message_bus_service/internal/app/weixin_open_event_consumer/worker" "zk_message_bus_service/internal/pkg/component" "zk_message_bus_service/internal/pkg/define" "zk_message_bus_service/internal/pkg/initialize" "zk_message_bus_service/pkg/consumer_manager" // 添加这个导入 "gitee.com/Sxiaobai/gs/gstool" ) // 全局消费者管理器实例 var globalConsumerManager *consumer_manager.Manager func Init() { initBase() initialize.InitLog(component.Env.LogPath, ``, 7) initialize.InitViper() initialize.InitHelper() initialize.InitMysql() initialize.InitRedis() initialize.InitModel() initialize.InitCache() initialize.InitOss() initialize.InitService() initialize.InitNsqPush() // 初始化消费者管理器(替换原来的NSQ初始化) initConsumerManager() component.Log.Debugf(`启动消息类消费者完成`) } // initConsumerManager 初始化消费者管理器 func initConsumerManager() { // 手动初始化NSQ结构体(仅生产者部分) lookUpHost := component.Viper.GetString(`nsq_look_up_host`) pubMsgHost := component.Viper.GetString(`nsq_pub_msg_host`) topic := component.Viper.GetString(`nsq_topic_event`) // 初始化NSQ结构体但不启动消费者 component.NsqEvent = initialize.QuickInitNsqProducerOnly(lookUpHost, pubMsgHost, topic, topic) // 创建配置 config := consumer_manager.DefaultConfig() // 创建NSQ适配器 nsqAdapter := &consumer_manager.NSQAdapter{ NsqStruct: component.NsqEvent, } // 创建管理器 var err error globalConsumerManager, err = consumer_manager.NewManager(config, nsqAdapter, component.MysqlZk, component.Log) if err != nil { component.Log.Errof("创建消费者管理器失败: %v", err) return } // 设置回调函数 globalConsumerManager.SetCallFunc(func(msg string, attempts uint16) bool { return worker.Work(msg) }) // 启动管理器 if err := globalConsumerManager.Start(); err != nil { component.Log.Errof("启动消费者管理器失败: %v", err) return } // 启动HTTP API服务器 if err := globalConsumerManager.StartAPIServer(); err != nil { component.Log.Errof("启动HTTP API服务器失败: %v", err) } // 初始化默认数量的消费者 if err := globalConsumerManager.AdjustCount(component.Env.ConsumerNums); err != nil { component.Log.Errof("初始化消费者失败: %v", err) } else { component.Log.Infof("消费者管理器启动成功,初始消费者数量: %d", component.Env.ConsumerNums) } } func Stop() { // 1. 停止消费者管理器(包括API服务器和NSQ生产者和消费者) if globalConsumerManager != nil { if err := globalConsumerManager.Stop(); err != nil { component.Log.Errof("停止消费者管理器失败: %v", err) } } component.Log.Infof(`停止完成`) _ = component.Log.Close() } ``` ### 步骤2: 确保NSQ适配器存在 确保在 `pkg/consumer_manager` 包中有NSQAdapter的实现,如果没有需要创建。 ### 步骤3: 配置环境变量 在配置文件中添加以下配置项: ```env # 消费者管理器配置(docker-compose.yml) environment: - TZ=Asia/Shanghai - SCALING_HOST_IP=x.x.x.x # 宿主机内网IP - SCALING_HOST_PORT=20111 # 宿主机映射端口 - SCALING_API_PORT=18888 # 固定监听端口 - SCALING_NSQ_LOOKUP_HOST=x.x.x.x:20061 # NSQ Lookup地址 - SCALING_NSQ_PUB_MSG_HOST=x.x.x.x:20060 # NSQ 发布地址 - SCALING_NSQ_HTTP_ADDRESS=x.x.x.x:20051 # NSQ HTTP地址 - SCALING_NSQ_ADMIN_ADDRESS=x.x.x.x:20071 # NSQ admin地址 - SCALING_NSQ_TOPIC=topic_name # NSQ Topic - SCALING_NSQ_CHANNEL=channel_name # NSQ Channel - SCALING_MAX_CONSUMERS=10 # 最大消费者数量 - SCALING_CONSUMER_NAME=xxxx_consumer #消费者名称 - SCALING_CONSUMER_REMARK=这个是干啥的消费者 #消费者中文备注 - SCALING_SPAWN_CONSUMER_COUNT=2 # 启动消费者数量 - SCALING_DEPTH_LIMIT=10000 # 消息积压报警阈值 ``` ### 步骤4: 验证集成 启动服务后,可以通过以下API验证集成是否成功(端口号为docker宿主机端口): ```bash # 检查健康状态 curl http://localhost:18087/api/v1/consumer/health # 查看消费者状态 curl http://localhost:18087/api/v1/consumer/status # 调整消费者数量 curl -X POST http://localhost:18087/api/v1/consumer/adjust \ -H "Content-Type: application/json" \ -d '{"consumer_count": 3}' # 暂停所有消费者 curl -X POST http://localhost:18087/api/v1/consumer/pause # 切换scaling_mode状态 curl -X POST http://localhost:18087/api/v1/consumer/toggle-scaling-mode ``` ### 步骤5: 数据库表创建 表结构如下: ```sql CREATE TABLE IF NOT EXISTS tbl_consumer_monitor ( id int(11) NOT NULL AUTO_INCREMENT COMMENT '自增ID', consumer_name varchar(100) NOT NULL COMMENT '消费者名称', server_ip varchar(50) NOT NULL COMMENT '服务器内网IP', port varchar(10) NOT NULL COMMENT '端口号', nsq_lookup_host varchar(100) NOT NULL COMMENT 'NSQ Lookup地址', nsq_http_address varchar(255) DEFAULT NULL COMMENT 'NSQ HTTP地址', nsq_admin_address varchar(255) DEFAULT NULL COMMENT 'NSQ Admin地址', topic varchar(100) NOT NULL COMMENT 'NSQ Topic', channel varchar(100) NOT NULL COMMENT 'NSQ Channel', max_consumer_count int(11) NOT NULL DEFAULT '10' COMMENT '最大消费者个数', consumer_count int(11) NOT NULL DEFAULT '1' COMMENT '消费者个数', spawn_consumer_count int(11) NOT NULL DEFAULT '1' COMMENT '启动消费者数量', depth_limit int(11) NOT NULL DEFAULT '10000' COMMENT '消息积压报警阈值', scaling_mode varchar(20) NOT NULL DEFAULT 'manual' COMMENT '扩缩容模式:auto/manual', create_time int(11) NOT NULL COMMENT '创建时间', update_time int(11) NOT NULL COMMENT '更新时间', consumer_remark varchar(100) NOT NULL DEFAULT '' COMMENT '消费者中文备注', PRIMARY KEY (id), UNIQUE KEY uk_consumer_server_port_topic_channel (consumer_name, server_ip, port, topic, channel) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='消费者监控表'; ``` ## API接口 ### 调整消费者数量 **接口**: `POST /api/v1/consumer/adjust` **功能**: 动态调整NSQ消费者数量 **请求参数**: ```json { "consumer_count": 3 } ``` **响应示例**: ```json { "code": 200, "message": "调整消费者数量成功", "data": { "current_count": 3, "max_consumer_count": 10, "spawn_consumer_count": 2, "depth_limit": 10000, "scaling_mode": "manual", "consumer_name": "tag_service", "server_ip": "192.168.1.100", "port": "18087" } } ``` ### 暂停消费者 **接口**: `POST /api/v1/consumer/pause` **功能**: 暂停所有消费者,将运行的消费者数量降至0 **请求参数**: 无需参数 **响应示例**: ```json { "code": 200, "message": "消费者暂停成功", "data": { "current_count": 0, "max_consumer_count": 10, "spawn_consumer_count": 2, "depth_limit": 10000, "scaling_mode": "manual", "consumer_name": "tag_service", "server_ip": "192.168.1.100", "port": "18087" } } ``` **使用场景**: - 系统维护时需要暂停消息处理 - 紧急情况下快速停止所有消费者 - 配合监控系统进行流量控制 ### 获取消费者状态 **接口**: `GET /api/v1/consumer/status` **功能**: 获取当前消费者运行状态 **响应示例**: ```json { "code": 200, "message": "获取消费者状态成功", "data": [ { "current_count": 3, "max_consumer_count": 10, "spawn_consumer_count": 2, "depth_limit": 10000, "scaling_mode": "manual", "consumer_name": "tag_service", "server_ip": "192.168.1.100", "port": "18087", "nsq_lookup_host": "nsqlookupd:4161", "nsq_http_address": "nsqd:4151", "nsq_admin_address": "nsqadmin:4171", "topic": "tag_events", "channel": "tag_service", "last_update": "2024-01-15 10:30:45", "consumer_remark": "标签服务消费者" } ] } ``` ### 切换扩缩容模式 **接口**: `POST /api/v1/consumer/toggle-scaling-mode` **功能**: 切换scaling_mode状态,在"auto"和"manual"之间切换 **请求参数**: 无需参数 **响应示例**: ```json { "code": 200, "message": "切换scaling_mode状态成功", "data": { "current_count": 3, "max_consumer_count": 10, "spawn_consumer_count": 2, "depth_limit": 10000, "scaling_mode": "auto", "consumer_name": "tag_service", "server_ip": "192.168.1.100", "port": "18087" } } ``` **使用场景**: - 在自动扩缩容和手动控制之间切换 - 根据业务需求动态调整扩缩容策略 - 配合监控系统进行智能化管理 ### 健康检查 **接口**: `GET /api/v1/consumer/health` **功能**: 检查服务健康状态 **响应示例**: ```json { "code": 200, "message": "服务正常", "data": { "status": "healthy", "time": "2024-01-15 10:30:45" } } ``` ## 新增字段说明 ### spawn_consumer_count - **类型**: int - **默认值**: 1 - **环境变量**: SCALING_SPAWN_CONSUMER_COUNT - **说明**: 启动消费者数量,当spawn_consumer_count大于max_consumer_count时,按照max_consumer_count进行上报 ### depth_limit - **类型**: int - **默认值**: 10000 - **环境变量**: SCALING_DEPTH_LIMIT - **说明**: 消息积压报警阈值,用于监控消息队列积压情况 ### scaling_mode - **类型**: string - **默认值**: "manual" - **可选值**: "auto", "manual" - **说明**: 扩缩容模式,auto为自动模式,manual为手动模式