一、项目背景
短链接服务是一个非常典型的后端综合项目:它既要处理高并发读写,又要兼顾数据一致性、可扩展性、可观测性与安全性。一个看似简单的“长链接转短链接”需求,在真实生产环境中会延伸出缓存、限流、认证、统计、异步解耦、容器化部署等一整套工程问题。
本文将使用 Gin + gRPC + MySQL + Redis + Kafka + Docker + K8s 实现一个分布式短链接服务,并围绕以下目标展开:
- 支持用户创建短链接
- 支持短链接解析与重定向
- 支持访问统计与异步数据分析
- 支持用户认证与权限隔离
- 支持高并发访问、缓存与限流
- 支持容器化部署与 Kubernetes 上线
为了让文章更具实战意义,文中给出的代码示例会尽量保持完整、可运行、可扩展。
二、需求分析与系统设计
2.1 功能需求
核心功能可以拆分为以下几类:
-
短链接生成
- 用户提交长链接
- 系统生成唯一短码
- 可设置过期时间
- 可设置自定义别名(可选)
-
短链接解析
- 用户访问短链接
- 系统查询原始长链接并返回 302/301 重定向
- 若链接不存在或已过期,返回错误页
-
访问统计
- 记录 PV、UV
- 记录访问时间、IP、UA、Referer
- 记录地域、设备、浏览器等信息
- 支持按日聚合分析
-
用户系统
- 用户注册、登录
- JWT 鉴权
- 用户只能管理自己的短链接
-
高并发能力
- 热点短链接缓存到 Redis
- 访问日志异步投递 Kafka
- 接口限流与系统保护
2.2 非功能需求
从工程角度,还需要关注以下指标:
- 高可用:服务实例可横向扩展
- 高性能:解析链路尽量走缓存,降低数据库压力
- 可观测:支持日志、监控、告警
- 可维护:服务职责清晰,便于拆分与扩展
- 安全性:鉴权、权限控制、防刷、防恶意跳转
2.3 系统架构设计
整体采用“接入层 + 核心服务层 + 存储与消息层”的分层方案:
flowchart LR
U[用户 / 浏览器] --> G[Gin API Gateway]
G --> A[Auth Service]
G --> S[ShortLink Service gRPC]
U --> R[Redirect Service Gin]
R --> C[Redis Cache]
R --> S
S --> M[(MySQL)]
R --> K[Kafka]
K --> T[Stats Consumer]
T --> M
T --> C
2.4 服务划分说明
| 服务 | 职责 |
|---|---|
| API Gateway(Gin) | 对外提供 REST API,处理用户请求、JWT 校验、参数校验 |
| ShortLink Service(gRPC) | 提供短链接创建、查询、删除、详情等核心能力 |
| Redirect Service(Gin) | 负责短链解析、缓存读取、重定向与访问日志投递 |
| Stats Consumer | 消费 Kafka 访问消息,异步写入统计库 |
| MySQL | 存储用户、短链接、访问统计聚合数据 |
| Redis | 缓存短链映射、热点数据、限流计数、幂等标记 |
| Kafka | 异步削峰解耦,处理访问日志与统计事件 |
2.5 核心设计思路
1)写请求走数据库,读请求优先走缓存
短链接创建时以 MySQL 为准,创建成功后写入 Redis。访问短链接时优先查询 Redis,如果缓存未命中,再回源 gRPC 服务查询 MySQL,并回填缓存。
2)统计链路异步化
用户访问短链接时,解析链路只做两件事:
- 返回跳转结果
- 将访问事件投递到 Kafka
统计计算、地域解析、浏览器分析等耗时逻辑统一在消费端异步处理,避免影响重定向延迟。
3)缓存与数据库双层设计
缓存适合处理高频读取,数据库适合存储最终一致的数据。短链接服务的经典模式是:
- Redis:处理高频解析
- MySQL:保证持久化与业务查询
4)限流与风控
短链接系统天然容易被刷,尤其是热门营销链接。因此必须在入口处做:
- IP 限流
- 用户维度限流
- 短码维度访问保护
- 非法 URL 检查
三、领域模型与数据库设计
3.1 领域模型
从业务角度,系统核心实体有四个:
- 用户(User)
- 短链接(ShortLink)
- 访问日志(VisitEvent)
- 访问统计聚合(LinkStatsDaily)
它们之间的关系如下:
| 实体 | 说明 |
|---|---|
| User | 系统用户,拥有多个短链接 |
| ShortLink | 短链接实体,包含短码、原始链接、状态、过期时间等 |
| VisitEvent | 每次访问产生一条事件日志,进入 Kafka 异步消费 |
| LinkStatsDaily | 某个短链接按天聚合后的统计数据 |
3.2 数据库表设计
3.2.1 用户表
CREATE TABLE users (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
username VARCHAR(64) NOT NULL UNIQUE,
password_hash VARCHAR(255) NOT NULL,
email VARCHAR(128) DEFAULT '',
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
3.2.2 短链接表
CREATE TABLE short_links (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
user_id BIGINT NOT NULL,
short_code VARCHAR(16) NOT NULL UNIQUE,
long_url TEXT NOT NULL,
status TINYINT NOT NULL DEFAULT 1,
expire_at DATETIME NULL,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_user_id (user_id),
INDEX idx_expire_at (expire_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
字段说明:
| 字段 | 含义 |
|---|---|
| short_code | 短码,例如 aZ91kd |
| long_url | 原始长链接 |
| status | 1 启用,0 禁用 |
| expire_at | 过期时间,可为空 |
3.2.3 访问事件表
如果访问量极大,访问明细通常不直接写 MySQL,而是先进入 Kafka,再根据需求落 ES、ClickHouse 或对象存储。为了演示完整链路,这里仍给出一份 MySQL 落表设计:
CREATE TABLE visit_events (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
short_code VARCHAR(16) NOT NULL,
user_id BIGINT NOT NULL DEFAULT 0,
ip VARCHAR(64) NOT NULL,
user_agent VARCHAR(255) NOT NULL,
referer VARCHAR(255) DEFAULT '',
country VARCHAR(64) DEFAULT '',
city VARCHAR(64) DEFAULT '',
device VARCHAR(64) DEFAULT '',
browser VARCHAR(64) DEFAULT '',
visited_at DATETIME NOT NULL,
INDEX idx_short_code_time (short_code, visited_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
3.2.4 按日统计表
CREATE TABLE link_stats_daily (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
short_code VARCHAR(16) NOT NULL,
stat_date DATE NOT NULL,
pv BIGINT NOT NULL DEFAULT 0,
uv BIGINT NOT NULL DEFAULT 0,
ip_count BIGINT NOT NULL DEFAULT 0,
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_short_code_date (short_code, stat_date)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
3.3 Redis Key 设计
Redis 不仅用于缓存短链映射,也用于限流和 UV 去重。
| Key | 示例 | 说明 |
|---|---|---|
shortlink:code:{code} |
shortlink:code:aZ91kd |
缓存短码到长链接映射 |
shortlink:uv:{code}:{date} |
shortlink:uv:aZ91kd:20260608 |
当日 UV 去重集合 |
rate_limit:create:{userID} |
rate_limit:create:1001 |
用户创建短链限流 |
rate_limit:redirect:{ip} |
rate_limit:redirect:10.0.0.1 |
IP 访问限流 |
3.4 Kafka Topic 设计
| Topic | 说明 |
|---|---|
shortlink.visit |
短链接访问事件 |
shortlink.deadletter |
异常消息死信队列 |
四、核心服务实现
下面给出一个精简但完整的示例,展示如何把多个组件串起来。
4.1 项目目录结构
distributed-shortlink/
├── cmd/
│ ├── api/main.go
│ ├── rpc/main.go
│ └── consumer/main.go
├── internal/
│ ├── config/config.go
│ ├── middleware/auth.go
│ ├── middleware/ratelimit.go
│ ├── model/model.go
│ ├── repo/repo.go
│ ├── service/shortlink.go
│ ├── handler/http.go
│ ├── grpcserver/server.go
│ └── mq/consumer.go
├── proto/shortlink.proto
├── go.mod
├── Dockerfile
├── docker-compose.yml
└── k8s/
├── mysql.yaml
├── redis.yaml
├── kafka.yaml
├── shortlink-api.yaml
└── shortlink-rpc.yaml
4.2 go.mod
module distributed-shortlink
go 1.23
require (
github.com/IBM/sarama v1.45.0
github.com/gin-gonic/gin v1.10.0
github.com/go-redis/redis/v8 v8.11.5
github.com/golang-jwt/jwt/v5 v5.2.1
github.com/google/uuid v1.6.0
github.com/joho/godotenv v1.5.1
golang.org/x/crypto v0.39.0
google.golang.org/grpc v1.74.2
google.golang.org/protobuf v1.36.6
gorm.io/driver/mysql v1.5.7
gorm.io/gorm v1.30.0
)
4.3 gRPC 协议定义
syntax = "proto3";
package shortlink;
option go_package = "distributed-shortlink/proto;pb";
service ShortLinkService {
rpc CreateShortLink(CreateShortLinkRequest) returns (CreateShortLinkResponse);
rpc ResolveShortLink(ResolveShortLinkRequest) returns (ResolveShortLinkResponse);
rpc GetShortLinkStats(GetShortLinkStatsRequest) returns (GetShortLinkStatsResponse);
}
message CreateShortLinkRequest {
int64 user_id = 1;
string long_url = 2;
int64 expire_at = 3;
}
message CreateShortLinkResponse {
string short_code = 1;
string short_url = 2;
}
message ResolveShortLinkRequest {
string short_code = 1;
}
message ResolveShortLinkResponse {
string long_url = 1;
bool valid = 2;
}
message GetShortLinkStatsRequest {
string short_code = 1;
}
message GetShortLinkStatsResponse {
string short_code = 1;
int64 pv = 2;
int64 uv = 3;
}
生成代码命令:
protoc --go_out=. --go-grpc_out=. proto/shortlink.proto
4.4 配置初始化
package config
import (
"fmt"
"os"
"github.com/IBM/sarama"
"github.com/go-redis/redis/v8"
"gorm.io/driver/mysql"
"gorm.io/gorm"
)
type AppConfig struct {
HTTPAddr string
GRPCAddr string
BaseURL string
JWTSecret string
RedisAddr string
KafkaBrokers []string
MySQLDSN string
}
func Load() AppConfig {
return AppConfig{
HTTPAddr: getenv("HTTP_ADDR", ":8080"),
GRPCAddr: getenv("GRPC_ADDR", ":9090"),
BaseURL: getenv("BASE_URL", "http://localhost:8080"),
JWTSecret: getenv("JWT_SECRET", "shortlink-secret"),
RedisAddr: getenv("REDIS_ADDR", "127.0.0.1:6379"),
KafkaBrokers: []string{getenv("KAFKA_ADDR", "127.0.0.1:9092")},
MySQLDSN: getenv("MYSQL_DSN", "root:root@tcp(127.0.0.1:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local"),
}
}
func NewDB(cfg AppConfig) (*gorm.DB, error) {
return gorm.Open(mysql.Open(cfg.MySQLDSN), &gorm.Config{})
}
func NewRedis(cfg AppConfig) *redis.Client {
return redis.NewClient(&redis.Options{Addr: cfg.RedisAddr})
}
func NewKafkaProducer(cfg AppConfig) (sarama.SyncProducer, error) {
sc := sarama.NewConfig()
sc.Producer.Return.Successes = true
sc.Producer.RequiredAcks = sarama.WaitForAll
sc.Producer.Retry.Max = 3
return sarama.NewSyncProducer(cfg.KafkaBrokers, sc)
}
func getenv(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
func Must(err error) {
if err != nil {
panic(fmt.Sprintf("fatal error: %v", err))
}
}
4.5 领域模型代码
package model
import "time"
type User struct {
ID int64 `gorm:"primaryKey"`
Username string `gorm:"size:64;uniqueIndex;not null"`
PasswordHash string `gorm:"size:255;not null"`
Email string `gorm:"size:128"`
CreatedAt time.Time
UpdatedAt time.Time
}
type ShortLink struct {
ID int64 `gorm:"primaryKey"`
UserID int64 `gorm:"index;not null"`
ShortCode string `gorm:"size:16;uniqueIndex;not null"`
LongURL string `gorm:"type:text;not null"`
Status int8 `gorm:"not null;default:1"`
ExpireAt *time.Time
CreatedAt time.Time
UpdatedAt time.Time
}
type LinkStatsDaily struct {
ID int64 `gorm:"primaryKey"`
ShortCode string `gorm:"size:16;not null;uniqueIndex:uk_code_date"`
StatDate time.Time `gorm:"type:date;not null;uniqueIndex:uk_code_date"`
PV int64 `gorm:"not null;default:0"`
UV int64 `gorm:"not null;default:0"`
IPCount int64 `gorm:"not null;default:0"`
CreatedAt time.Time
UpdatedAt time.Time
}
4.6 Repository 层
package repo
import (
"context"
"errors"
"time"
"distributed-shortlink/internal/model"
"gorm.io/gorm"
)
type Repository struct {
db *gorm.DB
}
func New(db *gorm.DB) *Repository {
return &Repository{db: db}
}
func (r *Repository) AutoMigrate() error {
return r.db.AutoMigrate(&model.User{}, &model.ShortLink{}, &model.LinkStatsDaily{})
}
func (r *Repository) CreateShortLink(ctx context.Context, link *model.ShortLink) error {
return r.db.WithContext(ctx).Create(link).Error
}
func (r *Repository) GetByShortCode(ctx context.Context, code string) (*model.ShortLink, error) {
var link model.ShortLink
err := r.db.WithContext(ctx).Where("short_code = ?", code).First(&link).Error
if err != nil {
return nil, err
}
if link.Status != 1 {
return nil, errors.New("short link disabled")
}
if link.ExpireAt != nil && link.ExpireAt.Before(time.Now()) {
return nil, errors.New("short link expired")
}
return &link, nil
}
func (r *Repository) UpsertDailyStats(ctx context.Context, code string, date time.Time, incPV, incUV int64) error {
sql := `
INSERT INTO link_stats_daily(short_code, stat_date, pv, uv, ip_count, created_at, updated_at)
VALUES(?, ?, ?, ?, ?, NOW(), NOW())
ON DUPLICATE KEY UPDATE
pv = pv + VALUES(pv),
uv = uv + VALUES(uv),
ip_count = ip_count + VALUES(ip_count),
updated_at = NOW();`
return r.db.WithContext(ctx).Exec(sql, code, date.Format("2006-01-02"), incPV, incUV, incUV).Error
}
func (r *Repository) GetStats(ctx context.Context, code string) (int64, int64, error) {
type result struct {
PV int64
UV int64
}
var res result
err := r.db.WithContext(ctx).
Table("link_stats_daily").
Select("COALESCE(SUM(pv),0) AS pv, COALESCE(SUM(uv),0) AS uv").
Where("short_code = ?", code).
Scan(&res).Error
return res.PV, res.UV, err
}
4.7 短链接生成算法
短码生成有很多方案,常见的有:
- 自增 ID + Base62 编码
- 雪花算法 + Base62 编码
- 随机字符串 + 唯一索引冲突重试
这里采用“随机字符串 + 唯一索引兜底”的实现,简单直接,适合中小型项目。
package service
import (
"context"
"crypto/rand"
"errors"
"fmt"
"net/url"
"time"
"distributed-shortlink/internal/model"
"distributed-shortlink/internal/repo"
"github.com/go-redis/redis/v8"
)
type ShortLinkService struct {
repo *repo.Repository
redis *redis.Client
baseURL string
}
func NewShortLinkService(r *repo.Repository, redis *redis.Client, baseURL string) *ShortLinkService {
return &ShortLinkService{repo: r, redis: redis, baseURL: baseURL}
}
func (s *ShortLinkService) Create(ctx context.Context, userID int64, longURL string, expireAt *time.Time) (string, string, error) {
if _, err := url.ParseRequestURI(longURL); err != nil {
return "", "", fmt.Errorf("invalid url: %w", err)
}
var code string
var err error
for i := 0; i < 5; i++ {
code, err = randomBase62(6)
if err != nil {
return "", "", err
}
link := &model.ShortLink{
UserID: userID,
ShortCode: code,
LongURL: longURL,
Status: 1,
ExpireAt: expireAt,
}
if err = s.repo.CreateShortLink(ctx, link); err == nil {
_ = s.redis.Set(ctx, cacheKey(code), longURL, 24*time.Hour).Err()
return code, fmt.Sprintf("%s/%s", s.baseURL, code), nil
}
}
return "", "", errors.New("generate short code failed after retries")
}
func (s *ShortLinkService) Resolve(ctx context.Context, code string) (string, error) {
if val, err := s.redis.Get(ctx, cacheKey(code)).Result(); err == nil && val != "" {
return val, nil
}
link, err := s.repo.GetByShortCode(ctx, code)
if err != nil {
return "", err
}
_ = s.redis.Set(ctx, cacheKey(code), link.LongURL, 24*time.Hour).Err()
return link.LongURL, nil
}
func (s *ShortLinkService) Stats(ctx context.Context, code string) (int64, int64, error) {
return s.repo.GetStats(ctx, code)
}
const alphabet = "0123456789abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ"
func randomBase62(n int) (string, error) {
b := make([]byte, n)
rb := make([]byte, n)
if _, err := rand.Read(rb); err != nil {
return "", err
}
for i := range b {
b[i] = alphabet[int(rb[i])%len(alphabet)]
}
return string(b), nil
}
func cacheKey(code string) string {
return "shortlink:code:" + code
}
4.8 gRPC 服务实现
package grpcserver
import (
"context"
"time"
"distributed-shortlink/internal/service"
pb "distributed-shortlink/proto"
)
type Server struct {
pb.UnimplementedShortLinkServiceServer
svc *service.ShortLinkService
}
func New(svc *service.ShortLinkService) *Server {
return &Server{svc: svc}
}
func (s *Server) CreateShortLink(ctx context.Context, req *pb.CreateShortLinkRequest) (*pb.CreateShortLinkResponse, error) {
var expireAt *time.Time
if req.ExpireAt > 0 {
t := time.Unix(req.ExpireAt, 0)
expireAt = &t
}
code, shortURL, err := s.svc.Create(ctx, req.UserId, req.LongUrl, expireAt)
if err != nil {
return nil, err
}
return &pb.CreateShortLinkResponse{ShortCode: code, ShortUrl: shortURL}, nil
}
func (s *Server) ResolveShortLink(ctx context.Context, req *pb.ResolveShortLinkRequest) (*pb.ResolveShortLinkResponse, error) {
longURL, err := s.svc.Resolve(ctx, req.ShortCode)
if err != nil {
return &pb.ResolveShortLinkResponse{Valid: false}, nil
}
return &pb.ResolveShortLinkResponse{LongUrl: longURL, Valid: true}, nil
}
func (s *Server) GetShortLinkStats(ctx context.Context, req *pb.GetShortLinkStatsRequest) (*pb.GetShortLinkStatsResponse, error) {
pv, uv, err := s.svc.Stats(ctx, req.ShortCode)
if err != nil {
return nil, err
}
return &pb.GetShortLinkStatsResponse{ShortCode: req.ShortCode, Pv: pv, Uv: uv}, nil
}
4.9 用户认证与权限管理
用户系统常见的实现方式是:
- 登录成功后签发 JWT
- Gin 中间件解析 JWT
- 将用户信息注入上下文
- 业务处理时校验资源归属
4.9.1 JWT 工具与认证中间件
package middleware
import (
"net/http"
"strings"
"time"
"github.com/gin-gonic/gin"
"github.com/golang-jwt/jwt/v5"
)
type Claims struct {
UserID int64 `json:"user_id"`
Username string `json:"username"`
jwt.RegisteredClaims
}
func GenerateToken(secret string, userID int64, username string) (string, error) {
claims := Claims{
UserID: userID,
Username: username,
RegisteredClaims: jwt.RegisteredClaims{
ExpiresAt: jwt.NewNumericDate(time.Now().Add(24 * time.Hour)),
IssuedAt: jwt.NewNumericDate(time.Now()),
},
}
token := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
return token.SignedString([]byte(secret))
}
func AuthMiddleware(secret string) gin.HandlerFunc {
return func(c *gin.Context) {
authHeader := c.GetHeader("Authorization")
if authHeader == "" || !strings.HasPrefix(authHeader, "Bearer ") {
c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "missing token"})
return
}
tokenStr := strings.TrimPrefix(authHeader, "Bearer ")
token, err := jwt.ParseWithClaims(tokenStr, &Claims{}, func(token *jwt.Token) (interface{}, error) {
return []byte(secret), nil
})
if err != nil || !token.Valid {
c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "invalid token"})
return
}
claims, ok := token.Claims.(*Claims)
if !ok {
c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "invalid claims"})
return
}
c.Set("userID", claims.UserID)
c.Set("username", claims.Username)
c.Next()
}
}
4.9.2 密码加密
package service
import "golang.org/x/crypto/bcrypt"
func HashPassword(password string) (string, error) {
b, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.DefaultCost)
return string(b), err
}
func VerifyPassword(hash string, password string) bool {
return bcrypt.CompareHashAndPassword([]byte(hash), []byte(password)) == nil
}
4.10 API Gateway 实现
下面的 Gin 网关负责:
- 创建短链接
- 获取统计信息
- 对外提供重定向入口
- 调用 gRPC 服务
package handler
import (
"context"
"encoding/json"
"net/http"
"time"
"github.com/IBM/sarama"
"github.com/gin-gonic/gin"
pb "distributed-shortlink/proto"
)
type HTTPHandler struct {
rpc pb.ShortLinkServiceClient
producer sarama.SyncProducer
}
func NewHTTPHandler(rpc pb.ShortLinkServiceClient, producer sarama.SyncProducer) *HTTPHandler {
return &HTTPHandler{rpc: rpc, producer: producer}
}
type CreateRequest struct {
LongURL string `json:"long_url" binding:"required,url"`
ExpireAt int64 `json:"expire_at"`
}
func (h *HTTPHandler) CreateShortLink(c *gin.Context) {
var req CreateRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
userID := c.GetInt64("userID")
resp, err := h.rpc.CreateShortLink(context.Background(), &pb.CreateShortLinkRequest{
UserId: userID,
LongUrl: req.LongURL,
ExpireAt: req.ExpireAt,
})
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{
"short_code": resp.ShortCode,
"short_url": resp.ShortUrl,
})
}
func (h *HTTPHandler) Redirect(c *gin.Context) {
code := c.Param("code")
resp, err := h.rpc.ResolveShortLink(context.Background(), &pb.ResolveShortLinkRequest{ShortCode: code})
if err != nil || !resp.Valid {
c.JSON(http.StatusNotFound, gin.H{"error": "short link not found"})
return
}
event := map[string]any{
"short_code": code,
"ip": c.ClientIP(),
"ua": c.Request.UserAgent(),
"referer": c.Request.Referer(),
"visited_at": time.Now().Unix(),
}
raw, _ := json.Marshal(event)
_, _, _ = h.producer.SendMessage(&sarama.ProducerMessage{
Topic: "shortlink.visit",
Value: sarama.ByteEncoder(raw),
})
c.Redirect(http.StatusFound, resp.LongUrl)
}
func (h *HTTPHandler) GetStats(c *gin.Context) {
code := c.Param("code")
resp, err := h.rpc.GetShortLinkStats(context.Background(), &pb.GetShortLinkStatsRequest{ShortCode: code})
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, gin.H{
"short_code": resp.ShortCode,
"pv": resp.Pv,
"uv": resp.Uv,
})
}
4.11 高并发处理与限流
4.11.1 基于 Redis 的滑动窗口限流
在生产环境中,创建接口和重定向接口都建议做限流。下面给出一个简单的固定窗口限流中间件。
package middleware
import (
"context"
"net/http"
"time"
"github.com/gin-gonic/gin"
"github.com/go-redis/redis/v8"
)
func RateLimitMiddleware(rdb *redis.Client, prefix string, limit int64, window time.Duration) gin.HandlerFunc {
return func(c *gin.Context) {
ctx := context.Background()
key := prefix + ":" + c.ClientIP()
cnt, err := rdb.Incr(ctx, key).Result()
if err != nil {
c.AbortWithStatusJSON(http.StatusInternalServerError, gin.H{"error": "rate limiter unavailable"})
return
}
if cnt == 1 {
_ = rdb.Expire(ctx, key, window).Err()
}
if cnt > limit {
c.AbortWithStatusJSON(http.StatusTooManyRequests, gin.H{"error": "too many requests"})
return
}
c.Next()
}
}
4.11.2 热点短链缓存
高并发场景下,绝大多数请求都集中在少数热门短链上。为此可以采用以下策略:
- 热点 key 提前预热到 Redis
- 延长热点 key 的过期时间
- 对不存在的短码做空值缓存,防止缓存穿透
- 在极端场景下对单个热点短码增加本地缓存
空值缓存示例:
func (s *ShortLinkService) Resolve(ctx context.Context, code string) (string, error) {
val, err := s.redis.Get(ctx, cacheKey(code)).Result()
if err == nil {
if val == "__nil__" {
return "", errors.New("not found")
}
return val, nil
}
link, err := s.repo.GetByShortCode(ctx, code)
if err != nil {
_ = s.redis.Set(ctx, cacheKey(code), "__nil__", 5*time.Minute).Err()
return "", err
}
_ = s.redis.Set(ctx, cacheKey(code), link.LongURL, 24*time.Hour).Err()
return link.LongURL, nil
}
4.12 访问统计与数据分析
访问日志不应该在重定向接口中同步落库,而应异步消费。下面给出一个 Kafka Consumer 示例。
package mq
import (
"context"
"encoding/json"
"log"
"time"
"distributed-shortlink/internal/repo"
"github.com/IBM/sarama"
"github.com/go-redis/redis/v8"
)
type VisitEvent struct {
ShortCode string `json:"short_code"`
IP string `json:"ip"`
UA string `json:"ua"`
Referer string `json:"referer"`
VisitedAt int64 `json:"visited_at"`
}
type Consumer struct {
repo *repo.Repository
redis *redis.Client
}
func NewConsumer(repo *repo.Repository, redis *redis.Client) *Consumer {
return &Consumer{repo: repo, redis: redis}
}
func (c *Consumer) Handle(msg *sarama.ConsumerMessage) {
var event VisitEvent
if err := json.Unmarshal(msg.Value, &event); err != nil {
log.Printf("invalid event: %v", err)
return
}
visitTime := time.Unix(event.VisitedAt, 0)
dateKey := visitTime.Format("20060102")
uvKey := "shortlink:uv:" + event.ShortCode + ":" + dateKey
ctx := context.Background()
isNewUV, err := c.redis.SAdd(ctx, uvKey, event.IP).Result()
if err == nil {
_ = c.redis.Expire(ctx, uvKey, 48*time.Hour).Err()
}
incUV := int64(0)
if isNewUV > 0 {
incUV = 1
}
if err := c.repo.UpsertDailyStats(ctx, event.ShortCode, visitTime, 1, incUV); err != nil {
log.Printf("upsert stats error: %v", err)
}
}
如果需要做进一步的数据分析,还可以在消费端增加以下维度:
| 维度 | 实现方式 |
|---|---|
| 地域分布 | 根据 IP 做 GeoIP 解析 |
| 浏览器分布 | 解析 User-Agent |
| 设备分布 | 区分移动端、PC、平板 |
| 来源分析 | 统计 Referer 域名 |
| 小时级趋势 | 按小时聚合访问量 |
4.13 服务启动代码
4.13.1 RPC 服务启动
package main
import (
"log"
"net"
"distributed-shortlink/internal/config"
"distributed-shortlink/internal/grpcserver"
"distributed-shortlink/internal/repo"
"distributed-shortlink/internal/service"
pb "distributed-shortlink/proto"
"google.golang.org/grpc"
)
func main() {
cfg := config.Load()
db, err := config.NewDB(cfg)
config.Must(err)
r := repo.New(db)
config.Must(r.AutoMigrate())
redisClient := config.NewRedis(cfg)
svc := service.NewShortLinkService(r, redisClient, cfg.BaseURL)
lis, err := net.Listen("tcp", cfg.GRPCAddr)
config.Must(err)
server := grpc.NewServer()
pb.RegisterShortLinkServiceServer(server, grpcserver.New(svc))
log.Printf("grpc server listening on %s", cfg.GRPCAddr)
config.Must(server.Serve(lis))
}
4.13.2 API 服务启动
package main
import (
"log"
"distributed-shortlink/internal/config"
"distributed-shortlink/internal/handler"
"distributed-shortlink/internal/middleware"
"github.com/gin-gonic/gin"
"google.golang.org/grpc"
pb "distributed-shortlink/proto"
)
func main() {
cfg := config.Load()
conn, err := grpc.Dial(cfg.GRPCAddr, grpc.WithInsecure())
config.Must(err)
rpcClient := pb.NewShortLinkServiceClient(conn)
producer, err := config.NewKafkaProducer(cfg)
config.Must(err)
h := handler.NewHTTPHandler(rpcClient, producer)
r := gin.Default()
r.Use(middleware.RateLimitMiddleware(config.NewRedis(cfg), "rate_limit:redirect", 300, 1))
auth := r.Group("/api")
auth.Use(middleware.AuthMiddleware(cfg.JWTSecret))
{
auth.POST("/links", middleware.RateLimitMiddleware(config.NewRedis(cfg), "rate_limit:create", 20, 60), h.CreateShortLink)
auth.GET("/links/:code/stats", h.GetStats)
}
r.GET("/:code", h.Redirect)
log.Printf("http server listening on %s", cfg.HTTPAddr)
config.Must(r.Run(cfg.HTTPAddr))
}
上述示例中的
RateLimitMiddleware第四个参数类型是time.Duration,因此生产代码中应写成1*time.Second、60*time.Second这样的形式。为方便阅读,本文在展示思路时省略了部分样板细节,实际落地时请按 Go 类型严格编写。
更严谨的写法如下:
auth.POST("/links",
middleware.RateLimitMiddleware(config.NewRedis(cfg), "rate_limit:create", 20, 60*time.Second),
h.CreateShortLink,
)
五、性能测试与优化
5.1 性能测试目标
短链接系统的关键性能指标主要集中在解析链路:
- 平均响应时间
- P95 / P99 延迟
- QPS
- 缓存命中率
- Kafka 堆积情况
- MySQL 慢查询数
5.2 使用 hey 进行压测
安装工具:
go install github.com/rakyll/hey@latest
压测重定向接口:
hey -n 10000 -c 200 http://127.0.0.1:8080/abc123
压测创建接口:
hey -n 2000 -c 50 -m POST \
-H "Authorization: Bearer <your-token>" \
-H "Content-Type: application/json" \
-d '{"long_url":"https://example.com/article/10001"}' \
http://127.0.0.1:8080/api/links
5.3 典型瓶颈分析
1)数据库成为瓶颈
若所有解析请求都回源 MySQL,会很快打满数据库连接池。解决方式:
- Redis 缓存短链映射
- 增加缓存 TTL
- 热点 key 预热
- 开启连接池与只读副本
2)统计写入拖慢主链路
如果重定向时同步写数据库,延迟会明显升高。解决方式:
- Kafka 异步写入
- 批量消费落库
- 统计与解析服务隔离部署
3)热点短链压力集中
某个爆款短链可能在短时间内涌入大量请求。可采用以下优化:
- Redis + 本地缓存双层缓存
- 访问限流
- 多副本部署 + K8s HPA 自动扩容
- 对热点 key 设置更长 TTL
5.4 优化建议清单
| 优化项 | 说明 |
|---|---|
| Redis 缓存 | 将解析链路从 DB 读转为缓存读 |
| Kafka 异步化 | 降低主请求耗时 |
| 连接池调优 | 优化 MySQL / Redis 连接数 |
| 批量写入 | 统计聚合批量落库减少 IO |
| 限流 | 避免恶意流量冲垮服务 |
| 空值缓存 | 防止缓存穿透 |
| 分库分表 | 超大规模数据场景下进一步扩展 |
5.5 MySQL 连接池优化示例
sqlDB, err := db.DB()
if err != nil {
panic(err)
}
sqlDB.SetMaxOpenConns(100)
sqlDB.SetMaxIdleConns(20)
sqlDB.SetConnMaxLifetime(30 * time.Minute)
5.6 Redis Pipeline 批量处理
对于统计类写操作,可以适当使用 Pipeline 提升吞吐:
pipe := redisClient.Pipeline()
pipe.Incr(ctx, "stats:pv:abc123")
pipe.SAdd(ctx, "stats:uv:abc123:20260608", ip)
_, err := pipe.Exec(ctx)
if err != nil {
panic(err)
}
六、部署上线与运维
6.1 Dockerfile
FROM golang:1.23 AS builder
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o shortlink-api ./cmd/api
RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o shortlink-rpc ./cmd/rpc
RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -o shortlink-consumer ./cmd/consumer
FROM alpine:3.20
WORKDIR /root/
COPY --from=builder /app/shortlink-api .
COPY --from=builder /app/shortlink-rpc .
COPY --from=builder /app/shortlink-consumer .
CMD ["./shortlink-api"]
6.2 本地联调 docker-compose.yml
version: '3.9'
services:
mysql:
image: mysql:8.0
environment:
MYSQL_ROOT_PASSWORD: root
MYSQL_DATABASE: shortlink
ports:
- "3306:3306"
command: --default-authentication-plugin=mysql_native_password
redis:
image: redis:7
ports:
- "6379:6379"
zookeeper:
image: bitnami/zookeeper:3.9
environment:
ALLOW_ANONYMOUS_LOGIN: yes
ports:
- "2181:2181"
kafka:
image: bitnami/kafka:3.7
ports:
- "9092:9092"
environment:
KAFKA_CFG_ZOOKEEPER_CONNECT: zookeeper:2181
ALLOW_PLAINTEXT_LISTENER: yes
KAFKA_CFG_LISTENERS: PLAINTEXT://:9092
KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
depends_on:
- zookeeper
rpc:
build: .
command: ["./shortlink-rpc"]
environment:
MYSQL_DSN: root:root@tcp(mysql:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local
REDIS_ADDR: redis:6379
KAFKA_ADDR: kafka:9092
BASE_URL: http://localhost:8080
GRPC_ADDR: :9090
depends_on:
- mysql
- redis
- kafka
api:
build: .
command: ["./shortlink-api"]
ports:
- "8080:8080"
environment:
MYSQL_DSN: root:root@tcp(mysql:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local
REDIS_ADDR: redis:6379
KAFKA_ADDR: kafka:9092
BASE_URL: http://localhost:8080
GRPC_ADDR: rpc:9090
HTTP_ADDR: :8080
JWT_SECRET: shortlink-secret
depends_on:
- rpc
consumer:
build: .
command: ["./shortlink-consumer"]
environment:
MYSQL_DSN: root:root@tcp(mysql:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local
REDIS_ADDR: redis:6379
KAFKA_ADDR: kafka:9092
depends_on:
- kafka
- mysql
- redis
启动命令:
docker compose up -d --build
6.3 Kubernetes 部署示例
6.3.1 API Deployment
apiVersion: apps/v1
kind: Deployment
metadata:
name: shortlink-api
spec:
replicas: 3
selector:
matchLabels:
app: shortlink-api
template:
metadata:
labels:
app: shortlink-api
spec:
containers:
- name: shortlink-api
image: shortlink:latest
command: ["./shortlink-api"]
ports:
- containerPort: 8080
env:
- name: MYSQL_DSN
value: root:root@tcp(mysql:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local
- name: REDIS_ADDR
value: redis:6379
- name: KAFKA_ADDR
value: kafka:9092
- name: BASE_URL
value: https://s.example.com
- name: GRPC_ADDR
value: shortlink-rpc:9090
- name: HTTP_ADDR
value: :8080
- name: JWT_SECRET
valueFrom:
secretKeyRef:
name: shortlink-secret
key: jwt-secret
---
apiVersion: v1
kind: Service
metadata:
name: shortlink-api
spec:
selector:
app: shortlink-api
ports:
- port: 8080
targetPort: 8080
6.3.2 RPC Deployment
apiVersion: apps/v1
kind: Deployment
metadata:
name: shortlink-rpc
spec:
replicas: 3
selector:
matchLabels:
app: shortlink-rpc
template:
metadata:
labels:
app: shortlink-rpc
spec:
containers:
- name: shortlink-rpc
image: shortlink:latest
command: ["./shortlink-rpc"]
ports:
- containerPort: 9090
env:
- name: MYSQL_DSN
value: root:root@tcp(mysql:3306)/shortlink?charset=utf8mb4&parseTime=True&loc=Local
- name: REDIS_ADDR
value: redis:6379
- name: BASE_URL
value: https://s.example.com
- name: GRPC_ADDR
value: :9090
---
apiVersion: v1
kind: Service
metadata:
name: shortlink-rpc
spec:
selector:
app: shortlink-rpc
ports:
- port: 9090
targetPort: 9090
6.3.3 HPA 自动扩容
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: shortlink-api-hpa
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: shortlink-api
minReplicas: 3
maxReplicas: 10
metrics:
- type: Resource
resource:
name: cpu
target:
type: Utilization
averageUtilization: 65
6.4 上线检查清单
| 检查项 | 说明 |
|---|---|
| 数据库索引 | 检查 short_code、user_id、统计维度索引 |
| Redis 持久化 | 根据业务选择 RDB / AOF 策略 |
| Kafka Topic 副本数 | 确保消息可靠性 |
| 日志采集 | 接入 ELK / Loki / Promtail |
| 指标监控 | 接入 Prometheus + Grafana |
| 告警规则 | 错误率、P99、Kafka Lag、MySQL 慢查询 |
| 灰度发布 | 新版本先小流量验证 |
6.5 运维建议
- 缓存命中率监控:如果短链解析缓存命中率下降,数据库压力会迅速升高。
- Kafka Lag 监控:一旦访问高峰期间消费者跟不上,统计会明显延迟。
- 链路追踪:建议接入 OpenTelemetry,排查慢请求更方便。
- 数据归档:访问明细数据量大,建议按时间分区或定期归档。
- 风控策略:对异常 IP、异常 Referer、恶意刷量行为增加黑名单机制。
七、总结
这个分布式短链接服务项目虽然业务模型简单,但非常适合作为 Golang 后端综合实战项目。通过它可以把常见的工程能力串成一套完整闭环:
- 使用 Gin 提供 REST 接口
- 使用 gRPC 拆分核心服务
- 使用 MySQL 存储业务核心数据
- 使用 Redis 提升解析性能并实现限流
- 使用 Kafka 解耦访问统计链路
- 使用 Docker 实现容器化交付
- 使用 Kubernetes 实现弹性扩缩容与上线部署
如果你只是做课程练习,可以先实现“创建短链接 + 访问跳转 + Redis 缓存”三部分;如果你要往生产级方案靠近,就要继续补齐认证、限流、异步统计、可观测性、灰度发布等能力。
从项目进阶路径来看,这个系统还可以继续演化:
- 自定义短码与冲突检测
- 链接标签、分组与搜索
- 管理后台
- 地域、设备、渠道等多维分析
- 多租户隔离
- 分库分表与热点治理
掌握了这个项目,基本也就掌握了一个真实 Go 微服务项目从设计、编码、性能优化到部署运维的完整过程。
📝 版权声明:本文为原创技术博客,转载请注明出处。
如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!