返回首页

Golang项目实践:4.8 综合项目之分布式短链接服务

一、项目背景

短链接服务是一个非常典型的后端综合项目:它既要处理高并发读写,又要兼顾数据一致性、可扩展性、可观测性与安全性。一个看似简单的“长链接转短链接”需求,在真实生产环境中会延伸出缓存、限流、认证、统计、异步解耦、容器化部署等一整套工程问题。

本文将使用 Gin + gRPC + MySQL + Redis + Kafka + Docker + K8s 实现一个分布式短链接服务,并围绕以下目标展开:

  • 支持用户创建短链接
  • 支持短链接解析与重定向
  • 支持访问统计与异步数据分析
  • 支持用户认证与权限隔离
  • 支持高并发访问、缓存与限流
  • 支持容器化部署与 Kubernetes 上线

为了让文章更具实战意义,文中给出的代码示例会尽量保持完整、可运行、可扩展

二、需求分析与系统设计

2.1 功能需求

核心功能可以拆分为以下几类:

  1. 短链接生成

    • 用户提交长链接
    • 系统生成唯一短码
    • 可设置过期时间
    • 可设置自定义别名(可选)
  2. 短链接解析

    • 用户访问短链接
    • 系统查询原始长链接并返回 302/301 重定向
    • 若链接不存在或已过期,返回错误页
  3. 访问统计

    • 记录 PV、UV
    • 记录访问时间、IP、UA、Referer
    • 记录地域、设备、浏览器等信息
    • 支持按日聚合分析
  4. 用户系统

    • 用户注册、登录
    • JWT 鉴权
    • 用户只能管理自己的短链接
  5. 高并发能力

    • 热点短链接缓存到 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.Second60*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_codeuser_id、统计维度索引
Redis 持久化 根据业务选择 RDB / AOF 策略
Kafka Topic 副本数 确保消息可靠性
日志采集 接入 ELK / Loki / Promtail
指标监控 接入 Prometheus + Grafana
告警规则 错误率、P99、Kafka Lag、MySQL 慢查询
灰度发布 新版本先小流量验证

6.5 运维建议

  1. 缓存命中率监控:如果短链解析缓存命中率下降,数据库压力会迅速升高。
  2. Kafka Lag 监控:一旦访问高峰期间消费者跟不上,统计会明显延迟。
  3. 链路追踪:建议接入 OpenTelemetry,排查慢请求更方便。
  4. 数据归档:访问明细数据量大,建议按时间分区或定期归档。
  5. 风控策略:对异常 IP、异常 Referer、恶意刷量行为增加黑名单机制。

七、总结

这个分布式短链接服务项目虽然业务模型简单,但非常适合作为 Golang 后端综合实战项目。通过它可以把常见的工程能力串成一套完整闭环:

  • 使用 Gin 提供 REST 接口
  • 使用 gRPC 拆分核心服务
  • 使用 MySQL 存储业务核心数据
  • 使用 Redis 提升解析性能并实现限流
  • 使用 Kafka 解耦访问统计链路
  • 使用 Docker 实现容器化交付
  • 使用 Kubernetes 实现弹性扩缩容与上线部署

如果你只是做课程练习,可以先实现“创建短链接 + 访问跳转 + Redis 缓存”三部分;如果你要往生产级方案靠近,就要继续补齐认证、限流、异步统计、可观测性、灰度发布等能力。

从项目进阶路径来看,这个系统还可以继续演化:

  • 自定义短码与冲突检测
  • 链接标签、分组与搜索
  • 管理后台
  • 地域、设备、渠道等多维分析
  • 多租户隔离
  • 分库分表与热点治理

掌握了这个项目,基本也就掌握了一个真实 Go 微服务项目从设计、编码、性能优化到部署运维的完整过程。


📝 版权声明:本文为原创技术博客,转载请注明出处。

如文章中存在错误或不准确之处,欢迎在评论区指正,感谢您的阅读与支持!

上一篇

Golang项目实践:4.5 消息队列与异步处理

下一篇

Golang项目实践:4.3 数据库与存储