首页
直播
壁纸
友链
搜索
1
微信小程序支付全链路实战:JSAPI 下单、调起支付、回调验签与退款
266 阅读
2
微信小程序云开发实战:云函数、云数据库与云存储的正确使用姿势
256 阅读
3
微信小程序自定义 tabBar 实战:custom-tab-bar 从适配到深色模式
255 阅读
4
微信小程序 Skyline 渲染引擎实战:worklet 动画从原理到落地
253 阅读
5
微信小程序分包进阶:独立分包、预下载与分包异步化实战
246 阅读
服务器运维
后端技术
前端技术
梯子
数据库
小程序
登录
搜索
标签搜索
fastadmin
Redis
微信小程序
前端开发
RabbitMQ
Go
服务器
codex
buildadmin
小程序
mysql
Nginx
Docker
Vue3
Node.js
MySQL优化
Linux
TypeScript
JWT
PHP
沿途的风景
累计撰写
74
篇文章
累计收到
0
条评论
首页
栏目
服务器运维
后端技术
前端技术
梯子
数据库
小程序
页面
直播
壁纸
友链
搜索到
24
篇与
» 后端技术
的结果
2026-03-09
Go + RabbitMQ 构建高并发任务队列实战
前言Go 语言的 goroutine 并发模型配合 RabbitMQ 的可靠消息投递,是构建高并发任务队列的经典方案。本文将实现一个完整的 Go + RabbitMQ 任务队列系统,包含生产者、消费者、重试机制和优雅退出。一、项目结构go-rabbitmq-queue/ ├── go.mod ├── config/ │ └── config.go ├── producer/ │ └── main.go ├── consumer/ │ └── main.go └── rabbitmq/ └── connection.go二、安装依赖go mod init github.com/yourname/go-rabbitmq-queue go get github.com/rabbitmq/amqp091-go三、RabbitMQ 连接封装package rabbitmq import ( "log" "time" amqp "github.com/rabbitmq/amqp091-go" ) type Config struct { URL string Exchange string Queue string RoutingKey string PrefetchCount int } type RabbitMQ struct { conn *amqp.Connection channel *amqp.Channel config Config } func New(cfg Config) (*RabbitMQ, error) { var conn *amqp.Connection var err error for i := 0; i < 5; i++ { conn, err = amqp.Dial(cfg.URL) if err == nil { break } log.Printf("连接失败(%d/5): %v", i+1, err) time.Sleep(3 * time.Second) } if err != nil { return nil, err } ch, err := conn.Channel() if err != nil { return nil, err } err = ch.ExchangeDeclare( cfg.Exchange, "direct", true, false, false, false, nil, ) if err != nil { return nil, err } args := amqp.Table{ "x-message-ttl": int32(60000), "x-dead-letter-exchange": cfg.Exchange + ".dlx", } _, err = ch.QueueDeclare( cfg.Queue, true, false, false, false, args, ) if err != nil { return nil, err } err = ch.QueueBind( cfg.Queue, cfg.RoutingKey, cfg.Exchange, false, nil, ) if err != nil { return nil, err } ch.Qos(cfg.PrefetchCount, 0, false) return &RabbitMQ{conn: conn, channel: ch, config: cfg}, nil } func (r *RabbitMQ) Channel() *amqp.Channel { return r.channel } func (r *RabbitMQ) Close() { r.channel.Close() r.conn.Close() }四、生产者package main import ( "encoding/json" "fmt" "log" "time" "github.com/yourname/go-rabbitmq-queue/rabbitmq" amqp "github.com/rabbitmq/amqp091-go" ) type Task struct { ID string `json:"id"` Type string `json:"type"` Payload interface{} `json:"payload"` } func main() { cfg := rabbitmq.Config{ URL: "amqp://admin:admin123@localhost:5672/", Exchange: "task.exchange", Queue: "task.queue", RoutingKey: "task.process", PrefetchCount: 10, } mq, err := rabbitmq.New(cfg) if err != nil { log.Fatal(err) } defer mq.Close() for i := 0; i < 100; i++ { task := Task{ ID: fmt.Sprintf("task-%d", i), Type: "email", Payload: map[string]string{ "to": fmt.Sprintf("user%d@example.com", i), "subject": "通知邮件", }, } body, _ := json.Marshal(task) err = mq.Channel().Publish( cfg.Exchange, cfg.RoutingKey, false, false, amqp.Publishing{ DeliveryMode: amqp.Persistent, ContentType: "application/json", Body: body, Timestamp: time.Now(), }, ) if err != nil { log.Printf("发送失败: %v", err) continue } log.Printf("发送任务: %s", task.ID) } log.Println("所有任务发送完成") }五、消费者(多 Worker 并发)package main import ( "encoding/json" "fmt" "log" "os" "os/signal" "sync" "syscall" "time" "github.com/yourname/go-rabbitmq-queue/rabbitmq" amqp "github.com/rabbitmq/amqp091-go" ) type Task struct { ID string `json:"id"` Type string `json:"type"` Payload interface{} `json:"payload"` } func processTask(task Task) error { log.Printf("处理任务: %s, 类型: %s", task.ID, task.Type) time.Sleep(500 * time.Millisecond) if time.Now().Unix()%10 == 0 { return fmt.Errorf("模拟处理失败") } log.Printf("任务完成: %s", task.ID) return nil } func startWorker(id int, mq *rabbitmq.RabbitMQ, wg *sync.WaitGroup) { defer wg.Done() msgs, err := mq.Channel().Consume( "task.queue", fmt.Sprintf("worker-%d", id), false, false, false, false, nil, ) if err != nil { log.Printf("Worker %d 启动失败: %v", id, err) return } log.Printf("Worker %d 启动", id) for msg := range msgs { var task Task if err := json.Unmarshal(msg.Body, &task); err != nil { log.Printf("Worker %d 解析失败: %v", id, err) msg.Nack(false, false) continue } if err := processTask(task); err != nil { log.Printf("Worker %d 处理失败: %s -> %v", id, task.ID, err) retryCount := getRetryCount(msg) if retryCount < 3 { msg.Nack(false, true) } else { log.Printf("Worker %d 任务 %s 重试超限,进入死信", id, task.ID) msg.Nack(false, false) } continue } msg.Ack(false) } log.Printf("Worker %d 退出", id) } func getRetryCount(msg amqp.Delivery) int { if deaths, ok := msg.Headers["x-death"].([]interface{}); ok && len(deaths) > 0 { if death, ok := deaths[0].(amqp.Table); ok { if count, ok := death["count"].(int64); ok { return int(count) } } } return 0 } func main() { cfg := rabbitmq.Config{ URL: "amqp://admin:admin123@localhost:5672/", Exchange: "task.exchange", Queue: "task.queue", RoutingKey: "task.process", PrefetchCount: 5, } mq, err := rabbitmq.New(cfg) if err != nil { log.Fatal(err) } defer mq.Close() var wg sync.WaitGroup workerCount := 5 for i := 1; i <= workerCount; i++ { wg.Add(1) go startWorker(i, mq, &wg) } sigs := make(chan os.Signal, 1) signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM) <-sigs log.Println("收到退出信号,等待 worker 完成...") mq.Close() wg.Wait() log.Println("所有 worker 已退出") }六、死信队列消费者func startDeadLetterConsumer(mq *rabbitmq.RabbitMQ) { ch := mq.Channel() ch.ExchangeDeclare("task.exchange.dlx", "fanout", true, false, false, false, nil) _, _ = ch.QueueDeclare("task.dlq", true, false, false, false, nil) _ = ch.QueueBind("task.dlq", "", "task.exchange.dlx", false, nil) msgs, _ := ch.Consume("task.dlq", "dlq-consumer", false, false, false, false, nil) go func() { for msg := range msgs { log.Printf("死信消息: %s", string(msg.Body)) msg.Ack(false) } }() }七、架构总结Producer → [task.exchange] → [task.queue] → 5个 Worker 并发消费 ↓ (失败/Nack) [task.exchange.dlx] → [task.dlq] → 死信消费者记录总结Go + RabbitMQ 的组合充分发挥了各自优势:Go 的 goroutine 让消费者可以轻松开几十个并发 worker,RabbitMQ 的 ACK 机制保证消息不丢失。生产环境注意:设置合理的 prefetch 避免消息堆积在单个 worker、实现死信队列处理失败消息、添加优雅退出逻辑避免消息中断。
2026年03月09日
10 阅读
0 评论
0 点赞
2026-03-02
Node.js JWT 认证完整实现:注册登录到权限控制
前言JWT(JSON Web Token)是现代 Web 应用中最流行的认证方案之一。本文将使用 Node.js + Express 实现一个完整的 JWT 认证系统,包含注册、登录、Token 刷新和权限控制。一、项目初始化mkdir jwt-auth && cd jwt-auth npm init -y npm install express jsonwebtoken bcryptjs cors dotenv npm install -D typescript @types/express @types/node ts-node-dev二、项目结构src/ ├── config/ │ └── index.ts ├── middleware/ │ └── auth.ts ├── routes/ │ └── auth.ts ├── utils/ │ └── jwt.ts ├── types/ │ └── index.ts └── app.ts三、核心实现3.1 JWT 工具类// src/utils/jwt.ts import jwt from "jsonwebtoken"; import { config } from "../config"; export interface JwtPayload { userId: string; role: string; } export function signAccessToken(payload: JwtPayload): string { return jwt.sign(payload, config.JWT_SECRET, { expiresIn: "15m", }); } export function signRefreshToken(payload: JwtPayload): string { return jwt.sign(payload, config.JWT_REFRESH_SECRET, { expiresIn: "7d", }); } export function verifyAccessToken(token: string): JwtPayload { return jwt.verify(token, config.JWT_SECRET) as JwtPayload; } export function verifyRefreshToken(token: string): JwtPayload { return jwt.verify(token, config.JWT_REFRESH_SECRET) as JwtPayload; }3.2 认证中间件// src/middleware/auth.ts import { Request, Response, NextFunction } from "express"; import { verifyAccessToken } from "../utils/jwt"; // 扩展 Request 类型 declare global { namespace Express { interface Request { user?: { userId: string; role: string }; } } } // 基本认证 export function authRequired(req: Request, res: Response, next: NextFunction) { const authHeader = req.headers.authorization; if (!authHeader?.startsWith("Bearer ")) { return res.status(401).json({ code: 401, msg: "未提供认证令牌" }); } const token = authHeader.substring(7); try { const payload = verifyAccessToken(token); req.user = { userId: payload.userId, role: payload.role }; next(); } catch { return res.status(401).json({ code: 401, msg: "令牌无效或已过期" }); } } // 角色权限控制 export function requireRole(...roles: string[]) { return (req: Request, res: Response, next: NextFunction) => { if (!req.user) { return res.status(401).json({ code: 401, msg: "未认证" }); } if (!roles.includes(req.user.role)) { return res.status(403).json({ code: 403, msg: "权限不足" }); } next(); }; }3.3 路由实现// src/routes/auth.ts import { Router, Response } from "express"; import bcrypt from "bcryptjs"; import { authRequired, requireRole } from "../middleware/auth"; import { signAccessToken, signRefreshToken, verifyRefreshToken } from "../utils/jwt"; const router = Router(); // 模拟用户存储(实际应使用数据库) const users: Map<string, { id: string; username: string; password: string; role: string }> = new Map(); const refreshTokens: Set<string> = new Set(); // 注册 router.post("/register", async (req, res: Response) => { const { username, password } = req.body; if (!username || !password) { return res.status(400).json({ code: 400, msg: "用户名和密码不能为空" }); } if (users.has(username)) { return res.status(409).json({ code: 409, msg: "用户已存在" }); } const hashedPassword = await bcrypt.hash(password, 10); const id = Date.now().toString(); users.set(username, { id, username, password: hashedPassword, role: "user" }); res.json({ code: 0, msg: "注册成功" }); }); // 登录 router.post("/login", async (req, res: Response) => { const { username, password } = req.body; const user = users.get(username); if (!user) { return res.status(404).json({ code: 404, msg: "用户不存在" }); } const valid = await bcrypt.compare(password, user.password); if (!valid) { return res.status(401).json({ code: 401, msg: "密码错误" }); } const payload = { userId: user.id, role: user.role }; const accessToken = signAccessToken(payload); const refreshToken = signRefreshToken(payload); refreshTokens.add(refreshToken); res.json({ code: 0, msg: "登录成功", data: { accessToken, refreshToken, user: { id: user.id, username: user.username, role: user.role } }, }); }); // 刷新 Token router.post("/refresh", (req, res: Response) => { const { refreshToken } = req.body; if (!refreshToken || !refreshTokens.has(refreshToken)) { return res.status(401).json({ code: 401, msg: "无效的刷新令牌" }); } try { const payload = verifyRefreshToken(refreshToken); refreshTokens.delete(refreshToken); const newPayload = { userId: payload.userId, role: payload.role }; const newAccessToken = signAccessToken(newPayload); const newRefreshToken = signRefreshToken(newPayload); refreshTokens.add(newRefreshToken); res.json({ code: 0, msg: "刷新成功", data: { accessToken: newAccessToken, refreshToken: newRefreshToken }, }); } catch { refreshTokens.delete(refreshToken); return res.status(401).json({ code: 401, msg: "刷新令牌已过期" }); } }); // 登出 router.post("/logout", (req, res: Response) => { const { refreshToken } = req.body; refreshTokens.delete(refreshToken); res.json({ code: 0, msg: "已登出" }); }); // 受保护的路由 router.get("/profile", authRequired, (req, res: Response) => { res.json({ code: 0, data: { userId: req.user!.userId, role: req.user!.role } }); }); // 管理员路由 router.get("/admin/users", authRequired, requireRole("admin"), (req, res: Response) => { const userList = Array.from(users.values()).map(u => ({ id: u.id, username: u.username, role: u.role, })); res.json({ code: 0, data: userList }); }); export default router;3.4 主应用// src/app.ts import express from "express"; import cors from "cors"; import dotenv from "dotenv"; import authRoutes from "./routes/auth"; dotenv.config(); const app = express(); app.use(cors()); app.use(express.json()); app.use("/api/auth", authRoutes); app.get("/", (_req, res) => { res.json({ msg: "JWT Auth API" }); }); app.listen(3000, () => { console.log("Server running on http://localhost:3000"); });四、安全最佳实践Access Token 短时效(15 分钟),Refresh Token 长时效(7 天)Refresh Token 一次性使用:每次刷新后旧 Token 失效HTTPS 传输:生产环境必须使用 HTTPS密码加密:使用 bcrypt,不要用 MD5/SHA1环境变量管理密钥:不要硬编码到代码中Token 黑名单:登出时将 Token 加入黑名单(Redis)总结JWT 认证方案无状态、易扩展,适合前后端分离架构。通过 Access Token + Refresh Token 双 Token 机制,可以在安全性和用户体验之间取得平衡。生产环境中建议将 Refresh Token 存储在 Redis 中,方便统一管理和吊销。
2026年03月02日
11 阅读
0 评论
0 点赞
2026-02-26
buildadmin框架执行 npm build 命令失败了
执行 npm build 命令失败了,您可以尝试删除 根目录/web/node_modules 和包管理器对应的 锁文件,锁文件名通常是:pnpm-lock.yamlpackage-lock.jsonyarn.lock推荐使用 pnpm 或 yarn NPM包管理器,安装 pnpm 命令:npm install -g pnpm手动执行 pnpm install 和 pnpm build 命令,成功后,再打开安装页面右下角的终端,点击 重新发布 按钮若执行失败,请升级 NodeJs 和 NPM若还是执行失败,请根据报错信息排除错误后,再继续进行安装。
2026年02月26日
12 阅读
0 评论
1 点赞
2026-02-04
thinkphp6 消息队列think-queue
1.安装队列依赖composer require topthink/think-queue2.配置文件/config/queue.php<?php return [ 'default' => 'redis', 'connections' => [ 'sync' => [ 'type' => 'sync', ], 'database' => [ 'type' => 'database', 'queue' => 'default', 'table' => 'jobs', ], 'redis' => [ 'type' => 'redis', 'queue' => 'default', 'host' => '127.0.0.1', 'port' => 6379, 'password' => '', 'select' => 0, 'timeout' => 0, 'persistent' => false, ], ], 'failed' => [ 'type' => 'none', 'table' => 'failed_jobs', ], ];3.在项目下新建一个Job目录存放处理消息4.控制器编写逻辑代码 app/controller/index.phpuse think\facade\Queue; public function job(Request $request) { $params = $request->get(); $jobHandlerClassName = 'app\job\Task'; $jobQueueName = 'task'; $orderData = ['order_sn'=>$params['id']]; //Queue::later();//立即执行 $isPushed = Queue::later(10, $jobHandlerClassName, $orderData, $jobQueueName); //这儿的10是指10秒后执行队列任务 if($isPushed !== false){ echo '队列添加成功'; }else{ echo '插入失败了'; } }5.编写对应的消费者类 app\job\task.php<?php namespace app\job; use think\queue\Job; class Task { public function fire(Job $job, $data) { $rt = $this->doJob($data); if($rt){ $job->delete(); return true; } // 重试三次失败 todo... if($job->attempts() == 3){ $job->delete(); return false; } //执行失败10S后重试 $job->release(10); } public function doJob($data) { echo date('Y-m-d H:i:s')."\n"; return false; } }php think queue:listen–queue helloJobQueue \ //监听的队列的名称–delay 0 \ //如果本次任务执行抛出异常且任务未被删除时,设置其下次执行前延迟多少秒,默认为0–memory 128 \ //该进程允许使用的内存上限,以 M 为单位–sleep 3 \ //如果队列中无任务,则多长时间后重新检查–tries 0 \ //如果任务已经超过重发次数上限,则进入失败处理逻辑,默认为0–timeout 60 // work 进程允许执行的最长时间,以秒为单位6.安装守护进程(Supervisor)
2026年02月04日
27 阅读
0 评论
0 点赞
2026-01-21
setEagerlyType字段理解
官方文档介绍:V5.0.4+版本开始一对一关联预载入支持两种方式:JOIN方式(一次查询)和IN方式(两次查询),如果要使用IN方式关联预载入,在关联定义方法中添加。 这句话的意思是jion方式关联和用in方式关联两种方式,in方式关联可以理解为先查出id,在用查出来的id去in另外一个表,也就是做了两次查询。 以下例子:当设置setEagerlyType为1时,由于用的是in方式,所以用category.name是查不到关联数据的 //在goods模型中设置关联模型category public function category(){ return $this->belongsTo('Category', 'goods_category_id', 'id', [], 'LEFT')->setEagerlyType(0); //return $this->belongsTo('category','goods_category_id','id'); } //goods表关联分类表 $list = $this->model->with(['category'])->where($where)->order($sort, $order) ->paginate($limit);
2026年01月21日
17 阅读
0 评论
0 点赞
1
...
3
4
5
0:00