文档目录

9.1 配套代码:短链服务骨架与六颗埋雷

对应小节:9.1 项目设定与埋雷清单 这是唯一一节需要写较多业务代码的地方——之后的章节都是编排与模板。

一、目录结构

src/main/kotlin/shortlink/
├── Application.kt              # Ktor 启动 + 指标注册
├── model/
│   ├── Link.kt                 # 领域模型
│   └── Dto.kt                  # request/response
├── service/
│   ├── LinkService.kt          # 业务逻辑(含埋雷 ⑥)
│   └── CodeGenerator.kt        # 短码生成(含埋雷 ⑥)
├── repository/
│   ├── LinkRepository.kt       # 数据访问(含埋雷 ① ② ④)
│   └── LinkCache.kt            # 缓存(含埋雷 ③)
├── http/
│   └── Routes.kt               # 路由(含埋雷 ⑤)
└── metrics/
    └── Metrics.kt              # 四层指标

二、领域模型与 DTO

// model/Link.kt
package shortlink.model

import java.time.Instant

data class Link(
    val id: Long,
    val code: String,
    val url: String,
    val userId: Long,
    val createdAt: Instant,
    val hits: Long,
)
// model/Dto.kt
package shortlink.model

import kotlinx.serialization.Serializable

@Serializable
data class CreateLinkRequest(val url: String, val userId: Long = 1)

@Serializable
data class CreateLinkResponse(val code: String)

三、埋雷 ⑥:synchronized 生成短码

// service/CodeGenerator.kt
package shortlink.service

import java.security.SecureRandom

private const val BASE62 = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"
private const val CODE_LENGTH = 6

class CodeGenerator {
    private val random = SecureRandom()

    /**
     * ═══ 埋雷 ⑥ ═══
     *
     * 问题:
     *   ① 用了 synchronized(线程级锁)—— 在协程里可能是"跨越挂起点的锁"
     *   ② 每次调用都新建 SecureRandom(其实可以用一个共享实例)
     *   ③ 冲突后的重试逻辑在 Service 层(见 LinkService)
     *
     * 症状:写路径偶发卡顿;线程 dump 可能出现 BLOCKED
     * 修法:改用 Mutex(协程感知),或去掉锁用 CAS 重试
     */
    @Synchronized
    fun generate(): String {
        val sb = StringBuilder(CODE_LENGTH)
        repeat(CODE_LENGTH) {
            sb.append(BASE62[random.nextInt(BASE62.length)])
        }
        return sb.toString()
    }
}

四、埋雷 ① ② ④:Repository 的三个问题

// repository/LinkRepository.kt
package shortlink.repository

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import shortlink.model.Link
import java.sql.Connection
import java.sql.ResultSet
import java.time.Instant
import javax.sql.DataSource

class LinkRepository(private val ds: DataSource) {

    /**
     * ═══ 埋雷 ① ═══
     * 建表时 code 列【没有唯一索引】→ 这条查询走全表扫描
     *
     * 症状:随数据量线性变慢(100 万行时约 380ms)
     * 验证:EXPLAIN (ANALYZE, BUFFERS) 会看到 Seq Scan
     * 修法:CREATE UNIQUE INDEX ON links(code)
     *
     * ═══ 埋雷 ④ ═══
     * 这里【没有】withContext(Dispatchers.IO)
     * → JDBC 是阻塞调用,直接跑在调用者的调度器上
     * → 如果调用者是 Dispatchers.Default(并行度=核数),会被占满
     *
     * 症状:CPU 低但所有接口一起慢
     * 验证:wall 火焰图 + RUNNABLE 线程数 == CPU 核数
     * 修法:withContext(Dispatchers.IO) + Semaphore 限流
     */
    suspend fun findByCode(code: String): Link? {
        return ds.connection.use { c ->              // ❌ 阻塞调用,未切调度器
            c.prepareStatement(
                "SELECT id, code, url, user_id, created_at, hits FROM links WHERE code = ?"
            ).use { ps ->
                ps.setString(1, code)
                ps.executeQuery().use { rs -> if (rs.next()) rs.toLink() else null }
            }
        }
    }

    suspend fun findById(id: Long): Link? {
        return ds.connection.use { c ->
            c.prepareStatement(
                "SELECT id, code, url, user_id, created_at, hits FROM links WHERE id = ?"
            ).use { ps ->
                ps.setLong(1, id)
                ps.executeQuery().use { rs -> if (rs.next()) rs.toLink() else null }
            }
        }
    }

    /**
     * ═══ 埋雷 ② ═══
     * 每次跳转都 UPDATE hits —— 写放大 + 行锁竞争
     *
     * 症状:高并发下热门短链的 UPDATE 争抢同一行锁;
     *       lock 火焰图有内容;pg_stat_activity 有 wait_event_type=Lock
     * 修法:批量/异步累加(但收益 < MDD,本 Lab 会否决它)
     */
    suspend fun incrementHits(id: Long) {
        ds.connection.use { c ->
            c.prepareStatement("UPDATE links SET hits = hits + 1 WHERE id = ?").use { ps ->
                ps.setLong(1, id)
                ps.executeUpdate()
            }
        }
    }

    suspend fun insert(code: String, url: String, userId: Long): Link {
        return ds.connection.use { c ->
            c.prepareStatement(
                "INSERT INTO links (code, url, user_id) VALUES (?, ?, ?) RETURNING id, created_at, hits"
            ).use { ps ->
                ps.setString(1, code)
                ps.setString(2, url)
                ps.setLong(3, userId)
                ps.executeQuery().use { rs ->
                    rs.next()
                    Link(rs.getLong(1), code, url, userId, rs.getTimestamp(2).toInstant(), rs.getLong(3))
                }
            }
        }
    }

    private fun ResultSet.toLink() = Link(
        getLong("id"), getString("code"), getString("url"),
        getLong("user_id"), getTimestamp("created_at").toInstant(), getLong("hits"),
    )
}

五、埋雷 ③:缓存 TTL 无抖动

// repository/LinkCache.kt
package shortlink.repository

import com.github.benmanes.caffeine.cache.Cache
import com.github.benmanes.caffeine.cache.Caffeine
import shortlink.model.Link
import java.time.Duration

class LinkCache {

    /**
     * ═══ 埋雷 ③ ═══
     * 所有 key 用【同一个固定 TTL】→ 同一批写入的 key 会同时失效
     *
     * 症状:数据库 QPS 出现周期性脉冲(每 10 分钟一次)
     *       —— 只有浸泡测试(≥1 小时)才能观察到
     * 验证:把 DB QPS 的脉冲时刻与 TTL(10min) 对齐
     * 修法:TTL 加随机抖动(±20%)
     */
    private val cache: Cache<String, Link> = Caffeine.newBuilder()
        .maximumSize(100_000)
        .expireAfterWrite(Duration.ofMinutes(10))       // ❌ 无抖动
        .recordStats()
        .build()

    fun get(code: String): Link? = cache.getIfPresent(code)

    fun put(code: String, link: Link) {
        cache.put(code, link)
    }

    fun invalidate(code: String) {
        cache.invalidate(code)
    }

    fun hitRate(): Double = cache.stats().hitRate()
    fun stats() = cache.stats()
}

六、Service 层:编排 + 短码冲突重试

// service/LinkService.kt
package shortlink.service

import shortlink.model.Link
import shortlink.repository.LinkCache
import shortlink.repository.LinkRepository
import org.slf4j.LoggerFactory

class LinkService(
    private val repo: LinkRepository,
    private val cache: LinkCache,
    private val codeGenerator: CodeGenerator,
) {
    private val log = LoggerFactory.getLogger(LinkService::class.java)

    /**
     * 读取路径:缓存优先,未命中查库并回填
     *
     * 注意:这里没有 withContext(Dispatchers.IO) —— 由 Repository 负责(而它也没有,
     *      所以整个链路跑在调用者的调度器上 → 埋雷 ④ 生效)
     */
    suspend fun resolve(code: String): Link? {
        // ① 先查缓存
        cache.get(code)?.let { return it }

        // ② 未命中 → 查库
        val link = repo.findByCode(code) ?: return null

        // ③ 回填缓存
        cache.put(code, link)

        // ④ 累加计数(埋雷 ②)
        repo.incrementHits(link.id)

        return link
    }

    /** 写入路径:生成短码 + 冲突重试 */
    suspend fun create(url: String, userId: Long): String {
        repeat(5) { attempt ->
            val code = codeGenerator.generate()        // 埋雷 ⑥:synchronized
            try {
                repo.insert(code, url, userId)
                return code
            } catch (e: Exception) {
                // 唯一索引冲突(code 已存在)—— 这是【正常业务】,不是故障
                if (isUniqueViolation(e)) {
                    log.debug("code collision on attempt {}: {}", attempt, code)
                    // 继续重试
                } else {
                    throw e
                }
            }
        }
        throw IllegalStateException("failed to generate unique code after 5 attempts")
    }

    private fun isUniqueViolation(e: Exception): Boolean =
        e.message?.contains("duplicate key") == true ||
        e.message?.contains("unique constraint") == true
}

注意 isUniqueViolation 这一段的含义:

短码冲突是【正常业务】——不是系统故障
→ 所以写接口的错误率阈值要区分两者(第 9.2 节)
→ 冲突后重试成功的不该计入错误率

七、埋雷 ⑤:热路径日志

// http/Routes.kt
package shortlink.http

import io.ktor.http.*
import io.ktor.server.application.*
import io.ktor.server.request.*
import io.ktor.server.response.*
import io.ktor.server.routing.*
import org.slf4j.LoggerFactory
import shortlink.model.*
import shortlink.service.LinkService

private val log = LoggerFactory.getLogger("shortlink.http")

fun Application.linkRoutes(service: LinkService) {
    routing {

        /**
         * ═══ 埋雷 ⑤ ═══
         * 每次跳转都打 INFO 日志,且包含了完整 URL 与 UA
         *
         * 症状:sys CPU 偏高(同步写磁盘)
         * 验证:JFR 的 jdk.FileWrite 事件;或做"降日志级别"的对照实验
         * 修法:异步 appender + 降级 + 参数化(不拼接大字符串)
         */
        get("/{code}") {
            val code = call.parameters["code"]!!

            // ❌ 热路径上的 INFO 日志(含完整字段)
            log.info(
                "redirect request: code={}, ua={}, ip={}",
                code,
                call.request.headers["User-Agent"] ?: "-",
                call.request.origin.remoteHost,
            )

            val link = service.resolve(code)
            if (link == null) {
                call.respond(HttpStatusCode.NotFound)
            } else {
                call.respondRedirect(link.url, permanent = false)
            }
        }

        post("/links") {
            val req = call.receive<CreateLinkRequest>()
            val code = service.create(req.url, req.userId)
            call.respond(HttpStatusCode.Created, CreateLinkResponse(code))
        }

        /** 对照接口:不查数据库、无日志 —— 用于验证"受影响的接口" */
        get("/health") {
            call.respondText("ok")
        }
    }
}

/health 这个对照接口很关键:

优化 2(调度器隔离)的效果要测 /health(第 9.7 节)
→ 因为它是"不查数据库但共用 Dispatchers.Default"的接口
→ 如果 Default 被占满,它也会慢(P99 从 8ms 涨到 380ms)

八、四层指标

// metrics/Metrics.kt
package shortlink.metrics

import com.github.benmanes.caffeine.cache.stats.CacheStats
import com.zaxxer.hikari.HikariDataSource
import io.micrometer.core.instrument.*
import io.micrometer.core.instrument.binder.jvm.*
import io.micrometer.core.instrument.binder.system.ProcessorMetrics
import io.micrometer.core.instrument.distribution.DistributionStatisticConfig
import io.micrometer.prometheusmetrics.PrometheusConfig
import io.micrometer.prometheusmetrics.PrometheusMeterRegistry
import java.time.Duration
import java.util.concurrent.atomic.AtomicInteger

val registry: PrometheusMeterRegistry = PrometheusMeterRegistry(PrometheusConfig.DEFAULT)

/** 活跃协程数(用于观察调度器压力) */
val activeCoroutines = AtomicInteger(0)

fun registerMetrics(ds: HikariDataSource, cacheStats: () -> CacheStats) {

    // ① 业务层
    Counter.builder("app.requests.total").register(registry)
    Counter.builder("app.requests.failed").register(registry)

    // ② 延迟层(直方图 + 围绕 SLO 的桶)
    Timer.builder("app.request.duration")
        .publishPercentileHistogram(true)
        .serviceLevelObjectives(
            Duration.ofMillis(1), Duration.ofMillis(2), Duration.ofMillis(5),
            Duration.ofMillis(10), Duration.ofMillis(25), Duration.ofMillis(50),
            Duration.ofMillis(75), Duration.ofMillis(100),   // ← SLO 附近加密
            Duration.ofMillis(150), Duration.ofMillis(250), Duration.ofMillis(500),
        )
        .register(registry)

    // ③ 资源层
    JvmMemoryMetrics().bindTo(registry)
    JvmGcMetrics().bindTo(registry)
    JvmThreadMetrics().bindTo(registry)
    ProcessorMetrics().bindTo(registry)

    // ④ 饱和度层 ⭐(本项目最重要的诊断指标)
    Gauge.builder("db.pool.pending") { ds.hikariPoolMXBean?.threadsAwaitingConnection ?: 0 }
        .description("等待连接的线程数(> 0 说明排队)")
        .register(registry)
    Gauge.builder("db.pool.active") { ds.hikariPoolMXBean?.activeConnections ?: 0 }
        .register(registry)
    Gauge.builder("db.pool.total") { ds.hikariPoolMXBean?.totalConnections ?: 0 }
        .register(registry)
    Gauge.builder("app.coroutines.active") { activeCoroutines.get().toDouble() }
        .description("活跃协程数(调度器压力的代理指标)")
        .register(registry)
    Gauge.builder("app.cache.hit_rate") { cacheStats().hitRate() }
        .description("缓存命中率(基线约 75%)")
        .register(registry)
    Gauge.builder("app.cache.size") { cacheStats().evictionCount().toDouble() }
        .register(registry)
}

九、验证六颗埋雷

#!/usr/bin/env bash
# tools/verify-traps.sh
set -uo pipefail

PG="${PG_CONN:-postgresql://app:app@localhost:5432/shortlink}"
SRC="${SRC_DIR:-src/main/kotlin/shortlink}"

echo "═══ 验证六颗埋雷 ═══"
echo

FAIL=0

# 埋雷 ①:code 无唯一索引
echo "① code 列索引"
IDX=$(psql -tAc "SELECT count(*) FROM pg_indexes WHERE tablename='links' AND indexdef LIKE '%(code)%';" "$PG" 2>/dev/null || echo "?")
if [ "$IDX" = "0" ]; then
  echo "   ✅ 无 code 索引(埋雷 ① 生效)"
elif [ "$IDX" = "?" ]; then
  echo "   ⚠️  无法查询数据库"
else
  echo "   ❌ code 已有索引 —— 埋雷 ① 失效!请 DROP INDEX"
  FAIL=$((FAIL+1))
fi
echo

# 埋雷 ②:每次跳转 UPDATE
echo "② UPDATE hits"
grep -rq "UPDATE links SET hits" "$SRC" 2>/dev/null \
  && echo "   ✅ 存在 UPDATE hits(埋雷 ② 生效)" \
  || { echo "   ❌ 未找到 UPDATE hits"; FAIL=$((FAIL+1)); }
echo

# 埋雷 ③:固定 TTL
echo "③ 缓存固定 TTL"
if grep -rq "expireAfterWrite" "$SRC" 2>/dev/null; then
  echo "   ✅ 使用 expireAfterWrite(固定 TTL,埋雷 ③ 生效)"
elif grep -rq "expireAfter" "$SRC" 2>/dev/null; then
  echo "   ⚠️  用了 expireAfter —— 确认是否有抖动(有抖动则埋雷 ③ 失效)"
else
  echo "   ❌ 未找到缓存 TTL 配置"; FAIL=$((FAIL+1))
fi
echo

# 埋雷 ④:Default 上有阻塞调用
echo "④ 阻塞调用未切调度器"
if grep -rq "withContext(Dispatchers.IO)" "$SRC" 2>/dev/null; then
  echo "   ❌ 已有 withContext(Dispatchers.IO) —— 埋雷 ④ 失效"
  FAIL=$((FAIL+1))
else
  echo "   ✅ 未见 withContext(Dispatchers.IO)(埋雷 ④ 生效)"
fi
echo

# 埋雷 ⑤:热路径日志
echo "⑤ 热路径 INFO 日志"
N=$(grep -rc "log.info" "$SRC" 2>/dev/null | awk -F: '{s+=$2} END{print s+0}')
if [ "$N" -gt 0 ]; then
  echo "   ✅ 存在 $N 处 log.info(埋雷 ⑤ 生效)"
else
  echo "   ❌ 未找到 log.info"; FAIL=$((FAIL+1))
fi
echo

# 埋雷 ⑥:synchronized
echo "⑥ synchronized"
if grep -rq "synchronized\|@Synchronized" "$SRC" 2>/dev/null; then
  echo "   ✅ 存在 synchronized(埋雷 ⑥ 生效)"
  grep -rn "synchronized\|@Synchronized" "$SRC" | head -3 | sed 's/^/      /'
else
  echo "   ❌ 未找到 synchronized"; FAIL=$((FAIL+1))
fi
echo

echo "═══════════════════════════════════════════════"
if [ "$FAIL" -eq 0 ]; then
  echo "✅ 六颗埋雷全部生效 —— 可以开始实验"
else
  echo "❌ 有 $FAIL 颗埋雷失效 —— 请先补回来"
  echo
  echo "为什么必须补:"
  echo "  埋雷失效 = 定位练习缺少素材 = 你会花时间找一个不存在的问题"
fi
echo "═══════════════════════════════════════════════"

十、动手改造

改动 观察什么
跑 verify-traps.sh 确认六颗埋雷都在(这是开始实验的前提)
故意提前加上 code 索引 埋雷 ① 失效,后面的容量曲线会"太好"——理解埋雷的必要性
故意加上 withContext(Dispatchers.IO) 埋雷 ④ 失效,/health 不会受影响——优化 2 就无从演示
把数据量从 100 万改成 1 万 埋雷 ① 的症状会消失(全表扫描 1 万行很快)——理解"规模决定结论"
给 /health 加日志 它会变慢——对照接口必须保持"干净"

十一、这段代码的局限

  • CodeGenerator 用 SecureRandom:它的初始化较慢(会读取系统熵),在高并发下可能成为瓶颈——但这是刻意的埋雷 ⑥。
  • isUniqueViolation 用字符串匹配判断异常:生产应该用 SQLState(23505)——这里是简化。
  • 没有实现真正的"埋雷 ③ 修复"(TTL 抖动):留到步骤六。
  • activeCoroutines 需要手工计数(在路由里 increment/decrement),不是自动的——省略了这部分代码以保持简洁。
  • 这个服务的错误处理很简化(比如没有处理 URL 格式校验、没有限流)——刻意如此,专注于性能。