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 格式校验、没有限流)——刻意如此,专注于性能。