
day 20 結尾寫的是「下一篇補上完整的 CRUD,insert、select、update、deleteWhere 4 組 DSL 怎麼寫,transaction { } 到底包了什麼、巢狀的時候算幾個交易、rollback 的邊界在哪裡,還有 day 19 那個 read-modify-write 的 lost update 問題」
這篇就照這個順序走,4 組 DSL 一組一組寫,每一組送出去的 SQL 都印出來對,update 回傳的那個筆數順便把 day 19 欠的 lost update 補掉
transaction 的部分從最基本的邊界開始,commit 跟 rollback 各發生在哪一行,巢狀的 transaction { } 算幾個交易,isolation level 跟 readOnly 這 2 個連線參數各管什麼、H2 認不認,再來是 suspendTransaction,它掛著 suspend 關鍵字,看起來像是非阻塞的那條路,這件事得用實測來回答,量完之後順著往下走,一條連線加一條 thread 怎麼互相卡死,40 個併發打滿之後不相干的路由要等多久
insert、select、update、deleteWhere 4 組 DSL,跟它們送出去的 SQLupdate 回傳的筆數就是併發的答案,day 19 欠的那個 lost updatetransaction { } 的邊界,commit 跟 rollback 各發生在哪一行readOnly 是連線上的 2 個參數,其中一個 H2 不理你suspendTransaction 掛著 suspend,但沒有換掉任何一條 threadHello, Ktor! 都要等 3 秒半這篇新增 2 個檔案,src/test/kotlin/com/cashwu/todo/TodoCrudTest.kt 放 CRUD,src/test/kotlin/com/cashwu/todo/TransactionTest.kt 放交易語意
day 20 在 TodoDatabaseTest.kt 裡寫的那個 private withDatabase helper 搬到 src/test/kotlin/com/cashwu/todo/TestApp.kt,3 個檔案共用,private 拿掉,簽章多一個 poolSize,TodoDatabaseTest.kt 裡原本那份刪掉
fun <T> withDatabase(
name: String,
seed: List<Todo> = defaultTodos,
poolSize: Int = 2,
block: (Database) -> T,
): T {
val source = todoDataSource(
DatabaseConfig(
url = "jdbc:h2:mem:$name",
driver = "org.h2.Driver",
user = "sa",
password = "",
poolSize = poolSize,
)
)
return source.use { block(it.connectAndSeed(seed)) }
}
跟 day 20 那版只差 DatabaseConfig 的 poolSize 從寫死的 2 換成接參數,搬過來的時候這一行最容易漏,漏了的話後面要開 8 條連線的測試會拿到 2 條,等一下會看到那是什麼下場。todoDataSource 照設定開一個 HikariCP 連線池、connectAndSeed(seed) 連上去建表塞種子資料、use 結束就關掉,2 個函式都在 TodoDatabase.kt。每個測試自己開一個名字不同的 in-memory 資料庫,預設池子還是 2 條
下面 DSL 的測試都放在 TodoCrudTest.kt 的 class 裡面
先看新增,Todos.insert { } 回傳的不是筆數而是一個 statement 物件,id 從它身上拿
@Test
fun `insert hands back the id the database generated`() {
withDatabase("insert") { database ->
val created = transaction(database) {
val statement = Todos.insert {
it[title] = "倒垃圾"
it[done] = false
it[createdAt] = FIXED_NOW.atOffset(ZoneOffset.UTC)
}
assertEquals(4, statement[Todos.id])
statement.resultedValues!!.single().toTodo()
}
assertEquals(Todo(4, "倒垃圾", false, FIXED_NOW), created)
}
}
statement[Todos.id] 是資料庫 AUTO_INCREMENT 發的那個 4,statement.resultedValues 則是整列回來,接上 day 20 那個 ResultRow.toTodo() 就是一個 Todo,POST 端點要回 201 加上完整物件的時候,會省掉一次 SELECT,FIXED_NOW 是 day 18 放進 TestApp.kt 的那個固定時間 2026-10-10T12:00:00Z
查詢,selectAll() 之後接 where、orderBy、limit,每一段各自回一個 query,最後被 map 觸發執行
@Test
fun `filtering ordering and limiting are three calls on the same query`() {
withDatabase("query") { database ->
transaction(database) {
val titles = Todos.selectAll()
.where { Todos.done eq false }
.orderBy(Todos.id to SortOrder.DESC)
.limit(1)
.map { it[Todos.title] }
assertEquals(listOf("寫 day 05 的文章"), titles)
assertEquals(3L, Todos.selectAll().count())
}
}
}
eq 這個字要注意,0.x 的 where 簽章是 where(predicate: SqlExpressionBuilder.() -> Op<Boolean>),lambda 帶著 receiver,eq 這個成員函式直接就在 scope 裡
1.x 把 receiver 拿掉了,變成 where(predicate: () -> Op<Boolean>),成員版的比較運算子再也碰不到,所以要 import org.jetbrains.exposed.v1.core.eq,少了這行編譯器回的是 Unresolved reference 'eq',不是 deprecation 警告,那些成員版雖然也標了 deprecation level ERROR,但你根本踩不到它們,0.x 的不用 import 也能跑,抄過來會編譯失敗
count() 回的是 Long 不是 Int,斷言要寫 3L,想要指定欄位就用 select 代替 selectAll
val row = Todos.select(Todos.id, Todos.title).where { Todos.id eq 1 }.single()
0.x 挑欄位用的 slice() 在 1.x 已經拿掉了,select(vararg) 是它的替代品,1.5.0 裡還留著一個同名的 slice,那是處理陣列欄位的函式,不是同一件事
更新跟刪除這 2 個回的是筆數
@Test
fun `update answers with the number of rows it changed`() {
withDatabase("update") { database ->
transaction(database) {
assertEquals(1, Todos.update({ Todos.id eq 2 }) { it[done] = true })
assertEquals(0, Todos.update({ Todos.id eq 999 }) { it[done] = true })
assertEquals(true, Todos.selectAll().where { Todos.id eq 2 }.single()[Todos.done])
}
}
}
update 的第 1 個參數是 where 條件,第 2 個是要改的欄位,2 個都是 lambda,deleteWhere 只有條件
@Test
fun `deleteWhere answers with the number of rows it removed`() {
withDatabase("delete") { database ->
transaction(database) {
assertEquals(1, Todos.deleteWhere { Todos.id eq 2 })
assertEquals(0, Todos.deleteWhere { Todos.id eq 999 })
assertEquals(listOf(1, 3), Todos.selectAll().map { it[Todos.id] })
}
}
}
改不到、刪不到都回 0 而不是丟例外,這個 0 等一下有用
想看 Exposed 到底送出去什麼,把 src/main/resources/logback.xml 的 Exposed logger 臨時開到 DEBUG
<logger name="Exposed" level="DEBUG"/>
上面那幾段跑出來是這樣
INSERT INTO TODOS (TITLE, DONE, CREATED_AT) VALUES ('倒垃圾', FALSE, '2026-10-10T12:00:00Z')
SELECT TODOS.ID, TODOS.TITLE, TODOS.DONE, TODOS.CREATED_AT FROM TODOS WHERE TODOS.ID = 2
SELECT TODOS.ID, TODOS.TITLE, TODOS.DONE, TODOS.CREATED_AT FROM TODOS WHERE TODOS.DONE = FALSE ORDER BY TODOS.ID DESC LIMIT 1
SELECT TODOS.ID, TODOS.TITLE FROM TODOS WHERE TODOS.ID = 1
SELECT COUNT(*) FROM TODOS
UPDATE TODOS SET DONE=TRUE WHERE TODOS.ID = 2
DELETE FROM TODOS WHERE TODOS.ID = 2
一句 Kotlin 對一句 SQL,沒有隱藏的 N+1、沒有 lazy loading 的驚喜,這是 DSL 這條路相對 DAO 那條路最實際的好處,你寫的東西跟資料庫收到的東西長得一樣
這一行 logger 看完就從 logback.xml 拿掉了
day 19 在「跟 Relix 的對照」那節留了一句話,原話是「Relix day 29 還提過一個這篇沒碰到的問題,update() 的「讀出來 → copy → 寫回去」是一段 read-modify-write,就算換成 ConcurrentHashMap 也還是有 lost update,這篇的 TodoRepository 現在只有讀跟新增,等更新端點補上去就會遇到,那時候答案會從 synchronized 換成資料庫的 transaction,剛好是 day 21 的題目」
拿一個具體情境來測,8 個 thread 同時想把第 2 筆待辦標成完成,而且每一條都要知道自己有沒有搶到,開 thread 的部分先包成 TodoCrudTest.kt 裡的一個 private helper
private fun race(count: Int, body: (CyclicBarrier) -> Unit) {
val barrier = CyclicBarrier(count)
val threads = (1..count).map { Thread { body(barrier) } }
threads.forEach { it.start() }
threads.forEach { it.join() }
}
CyclicBarrier 交給 body 自己決定要卡在哪一行,先看照直覺寫的版本,讀出來看一眼,還沒完成就寫回去,這個測試的斷言寫 8 是故意的,它要證明的是 lost update 真的會發生,不是我們想要的行為,留著它是把「這樣寫會壞」這件事綁在測試裡,跟下一個測試是一組對照,哪天有人把 WHERE 的條件拿掉,下一個測試會紅,這一個還是綠,2 個對著看才知道差在哪
@Test
fun `reading before writing lets eight threads claim the same todo`() {
withDatabase("claim-read", poolSize = 8) { database ->
val winners = AtomicInteger()
race(8) { barrier ->
transaction(database) {
val row = Todos.selectAll().where { Todos.id eq 2 }.single()
barrier.await()
if (!row[Todos.done]) {
Todos.update({ Todos.id eq 2 }) { it[done] = true }
winners.incrementAndGet()
}
}
}
assertEquals(8, winners.get())
}
}
barrier 卡在讀完之後、寫之前,8 條 thread 都讀到 done = false 才一起往下走,結果是 8 個都認為自己搶到了,這不是 H2 的問題,也不是 Exposed 的問題,READ COMMITTED 這個隔離等級本來就允許這件事,讀的那一刻沒有鎖,8 個 SELECT 看到的都是舊值
poolSize = 8 要真的生效,池子只有 2 條的話,拿到連線的 2 條 thread 在 barrier 等另外 6 條,那 6 條在 getConnection() 等前面的人還連線,30 秒後 HikariCP 丟 timeout、Exposed 重試 3 次,最後那 6 條帶著例外死掉,barrier 永遠湊不齊 8 個人,這個測試會掛在那裡不結束,這件事後面「一條連線加一條 thread 就能卡死」那節會正式碰到
把條件搬進 WHERE 就好了
@Test
fun `putting the condition in the where clause leaves exactly one winner`() {
withDatabase("claim-guard", poolSize = 8) { database ->
val winners = AtomicInteger()
race(8) { barrier ->
transaction(database) {
barrier.await()
val changed = Todos.update({ (Todos.id eq 2) and (Todos.done eq false) }) {
it[done] = true
}
if (changed == 1) winners.incrementAndGet()
}
}
assertEquals(1, winners.get())
}
}
送出去的 SQL 是
UPDATE TODOS SET DONE=TRUE WHERE (TODOS.ID = 2) AND (TODOS.DONE = FALSE)
8 條 thread 同時發這一句,第 1 個拿到那一列的鎖,改完 commit,剩下 7 個在鎖上等,等到之後資料庫重新比對 WHERE,DONE 已經是 TRUE,條件不成立,回 0,重跑幾輪都是 1 個 1、7 個 0
所以 update 那個回傳值就是答案,判斷條件跟寫入是同一句 SQL,中間沒有縫可以插隊,changed == 1 代表這條 thread 搶到了,changed == 0 代表別人先做了,day 19 用 synchronized 解的是同一個 JVM 裡的競爭,這裡的解法連跨機器都成立,因為仲裁的人變成資料庫
and 也要 import,org.jetbrains.exposed.v1.core.and,但它的情況跟 eq 不一樣。它在 0.x 就已經是 top-level 函式了,從來沒有成員版,只是套件名從 org.jetbrains.exposed.sql 搬到了 org.jetbrains.exposed.v1.core
transaction { } 做的事情用一句話講完,是從連線池借一條連線、把 autoCommit 關掉、跑你的 block、正常回來就 commit、丟例外就 rollback,最後把連線還回去
rollback 這件事測起來很直接,先在 TransactionTest.kt 放一個 private helper,這個檔案要塞很多筆只有標題不同的資料
private fun addTodo(title: String) = Todos.insert {
it[Todos.title] = title
it[done] = false
it[createdAt] = FIXED_NOW.atOffset(ZoneOffset.UTC)
}
它自己不開交易,靠呼叫端那個 transaction { } 提供,Exposed 從執行緒上找得到目前這個交易
@Test
fun `an exception takes the whole block back with it`() {
withDatabase("rollback") { database ->
assertFailsWith<IllegalStateException> {
transaction(database) {
addTodo("會不見")
assertEquals(4L, Todos.selectAll().count())
error("後面炸了")
}
}
assertEquals(3L, transaction(database) { Todos.selectAll().count() })
}
}
block 裡面看到 4 筆,block 外面看到 3 筆。中間那個 assertEquals(4L, ...) 是刻意的,它證明 INSERT 真的送到資料庫了,只是還沒 commit,同一個交易裡自己看得到自己寫的東西
Exposed 1.5.0 的 Transactions.kt 裡這段邏輯是這樣寫的
return try {
block().also {
if (shouldCommit) {
transaction.commit()
}
}
} catch (cause: SQLException) {
val currentStatement = transaction.currentStatement
errorBlock(currentStatement)
throw cause
} catch (cause: Throwable) {
if (shouldCommit) {
val currentStatement = transaction.currentStatement
errorBlock(currentStatement)
}
throw cause
}
第 2 個 catch 收的是 Throwable,所以不管你丟的是 SQLException、IllegalStateException 還是 OutOfMemoryError,都會走到 rollback 那條路,例外本身照樣往外丟,這裡 shouldCommit 是下一節的主角
順帶一個容易忽略的行為,SQLException 這一支還會重試,inTopLevelTransaction 裡面用一個 while (true) 把整個 block 包起來,預設 defaultMaxAttempts = 3,也就是資料庫層級的錯誤會讓整個 block 重跑,連第 1 次總共最多 3 次,block 裡面如果有寫檔案、送訊息這種資料庫管不到的副作用,重試的時候會做第二遍
transaction { } 裡面再包一個 transaction { },直覺會以為內層是一個獨立的交易,內層失敗就內層 rollback。預設不是這樣
下面幾個測試都要看最後表裡剩下什麼,所以 TransactionTest.kt 裡有一個把標題撈成 List 的 private helper
private fun titles(database: Database) =
transaction(database) { Todos.selectAll().map { it[Todos.title] } }
有了它,內層失敗之後表裡剩什麼一句就問得出來
@Test
fun `a nested transaction shares the outer transaction by default`() {
withDatabase("nested-off") { database ->
assertEquals(false, database.useNestedTransactions)
transaction(database) {
runCatching {
transaction(database) {
addTodo("內層")
error("內層炸了")
}
}
}
assertContains(titles(database), "內層")
}
}
內層丟了例外,被 runCatching 接住,外層正常結束,那筆「內層」還在表裡
原因就在上一節那個 shouldCommit,transaction 發現自己在別的交易裡面的時候,走的是這條
executeTransactionWithErrorHandling(
transaction,
shouldCommit = outer.db.useNestedTransactions
) {
transaction.statement()
}
useNestedTransactions 預設是 false,所以 shouldCommit 是 false,內層仍在外層交易裡,只是沒有獨立的 commit、rollback 或 savepoint,整包的 commit 時機還是外層 block 結束的那一刻,內層寫的東西當然跟著進去
要真的分開就得打開那個開關,它在 Exposed 自己的 DatabaseConfig 上,跟 day 20 那個我們自己寫的 DatabaseConfig 是同名的 2 個型別,所以測試裡用 import alias 分開
import org.jetbrains.exposed.v1.core.DatabaseConfig as ExposedDatabaseConfig
TransactionTest.kt 開頭有一個 private helper 專門開這種資料庫
private fun <T> withNestingDatabase(name: String, block: (Database) -> T): T {
val source = todoDataSource(
DatabaseConfig("jdbc:h2:mem:$name", "org.h2.Driver", "sa", "", 2)
)
return source.use {
val database = Database.connect(
it,
databaseConfig = ExposedDatabaseConfig { useNestedTransactions = true },
)
transaction(database) { SchemaUtils.create(Todos) }
block(database)
}
}
它不走 connectAndSeed,因為那個函式沒有帶 Exposed config 的入口,而且這個測試要的是一張空表
@Test
fun `savepoints make the inner block roll back on its own`() {
withNestingDatabase("nested-on") { database ->
transaction(database) {
runCatching {
transaction(database) {
addTodo("內層")
error("內層炸了")
}
}
addTodo("外層")
}
assertEquals(listOf("外層"), titles(database))
}
}
只剩「外層」,打開之後 Exposed 在內層進去之前下一個 savepoint,內層失敗就回到那個點,外層繼續走,這是 JDBC 的 Connection.setSavepoint() 而不是 SQL 語句,所以 DEBUG log 裡看不到它
要不要打開這個開關,看你怎麼寫 service,如果內層的 transaction { } 只是為了「借一條連線來跑幾句 SQL」,那預設值反而是對的,整條呼叫鏈共用一個交易,語意乾淨,真的需要「這一小段失敗不影響外面」的時候再打開,而且要知道 savepoint 在高併發下是有成本的
transaction() 除了 db 還收 2 個參數,一個是隔離等級、一個是唯讀,這 2 個測試也放在 TransactionTest.kt 的 class 裡面
@Test
fun `the isolation level comes from the connection and can be raised`() {
withDatabase("isolation") { database ->
transaction(database) {
assertEquals(Connection.TRANSACTION_READ_COMMITTED, transactionIsolation)
}
transaction(database, transactionIsolation = Connection.TRANSACTION_SERIALIZABLE) {
assertEquals(Connection.TRANSACTION_SERIALIZABLE, transactionIsolation)
}
}
}
預設是 TRANSACTION_READ_COMMITTED,也就是常數 2,前面那 8 條 thread 全部搶到就是這個等級的正常行為,換成 TRANSACTION_SERIALIZABLE 資料庫會改用更嚴的方式擋,代價是衝突的時候有人會拿到錯誤而不是拿到 0,程式要準備好重試,這裡點到為止,真的要調它之前得先確認資料庫的實作方式,PostgreSQL 跟 H2 的 SERIALIZABLE 做法就不一樣
readOnly 這個參數則有一個要注意的地方
@Test
fun `readOnly is a hint and H2 writes anyway`() {
withDatabase("readonly") { database ->
transaction(database, readOnly = true) {
assertEquals(true, readOnly)
addTodo("唯讀交易寫進去的")
}
assertContains(titles(database), "唯讀交易寫進去的")
}
}
交易自己說它是唯讀的,然後照樣寫進去了,沒有例外,readOnly 最後落到 JDBC 的 Connection.setReadOnly(),那個方法在規格上就是給 driver 的最佳化提示,擋不擋是 driver 自己決定,H2 選擇不擋,PostgreSQL 這邊會擋,SQLSTATE 25006 的訊息是 cannot execute INSERT in a read-only transaction
前面每一個 transaction { } 都是同步呼叫,Exposed 另外提供 suspendTransaction,名字掛著 suspend,很容易讀成「這個版本是非阻塞的」
它不是,先看它跟 transaction 差在哪裡,Transactions.kt 裡 2 個函式的簽章擺在一起就清楚了
fun <T> transaction(
db: Database? = null,
transactionIsolation: Int? = db?.transactionManager?.defaultIsolationLevel,
readOnly: Boolean? = db?.transactionManager?.defaultReadOnly,
statement: JdbcTransaction.() -> T
): T
suspend fun <T> suspendTransaction(
db: Database? = null,
transactionIsolation: Int? = db?.transactionManager?.defaultIsolationLevel,
readOnly: Boolean? = db?.transactionManager?.defaultReadOnly,
statement: suspend JdbcTransaction.() -> T
): T
差別在 statement 那一行,suspend JdbcTransaction.() -> T,它給你的能力是「在交易裡面呼叫別的 suspend 函式」,例如中間要打一個 HTTP API,至於 Todos.insert、Todos.selectAll 這些,執行的還是 exposed-jdbc 那套 PreparedStatement.executeQuery(),JDBC 的 API 從頭到尾沒有非阻塞版本,要走真正非阻塞的那條路,Exposed 1.x 另外有一個 exposed-r2dbc,那是換掉整個驅動層,不是這個模組的事
會不會換 dispatcher 呢,答案在 JdbcTransaction.kt 這個函式
suspend fun <T> withTransactionContext(transaction: JdbcTransaction, block: suspend CoroutineScope.() -> T): T {
val dispatcher = currentCoroutineContext()[CoroutineDispatcher.Key]
val context = if (dispatcher != null) {
transaction.asContext()
} else {
transaction.asContext() + transaction.db.config.dispatcher
}
return withContext(context, block)
}
只有在目前的 context 完全沒有 dispatcher 的時候,它才會補上 db.config.dispatcher,那個東西預設是 Dispatchers.IO,Ktor 的 handler 一定有 dispatcher,所以走的是上面那條,原地不動
直接在 Netty 上驗。在 Application.kt 的 routing { } 裡臨時加一個路由,把 handler 所在的 dispatcher、handler 自己的 thread、2 種交易裡面的 thread,4 個名字一起印出來
get("/thread") {
val outer = Thread.currentThread().name
val inTx = transaction(database) { Thread.currentThread().name }
val inSuspend = suspendTransaction(database) { Thread.currentThread().name }
call.respondText(
"""
dispatcher=${currentCoroutineContext()[CoroutineDispatcher]}
outer=$outer
inTx=$inTx
inSuspend=$inSuspend
""".trimIndent()
)
}
database 是 module() 裡 day 20 就有的那個 val database: Database by dependencies。dispatcher 要用 currentCoroutineContext() 拿,handler 的 receiver RoutingContext 自己也有一個 coroutineContext 屬性,從那個查 CoroutineDispatcher 會拿到 null
./gradlew run 起來打一次 GET /thread
dispatcher=NettyDispatcher@54e0d4d9
outer=eventLoopGroupProxy-4-2
inTx=eventLoopGroupProxy-4-2
inSuspend=eventLoopGroupProxy-4-2
eventLoopGroupProxy-4-2 是 Netty 的 event loop thread,handler 跑在它上面,transaction 裡面是它,suspendTransaction 裡面還是它,dispatcher 那行是 Netty 自己的 NettyDispatcher,對回上面 withTransactionContext 那段原始碼,有 dispatcher 就原地不動,這條路由測完就刪掉了
那「佔住 thread」到底有多具體 ? 讓 4 個交易在同一條 thread 上跑,每個做一件要 200 毫秒的事。先在 TransactionTest.kt 放一個 private helper onOneThread,開一條有名字的 thread 當 dispatcher,跑完關掉,後面卡死那個測試也會用它
private fun <T> onOneThread(name: String = "只有這條", block: suspend CoroutineScope.() -> T): T {
val pool = Executors.newSingleThreadExecutor { runnable -> Thread(runnable, name) }
try {
return runBlocking { withContext(pool.asCoroutineDispatcher(), block) }
} finally {
pool.shutdownNow()
}
}
有了它,測試本身就是 2 組各 4 個協程,一組跑 suspendTransaction,一組跑 delay
@Test
fun `four suspendTransaction blocks on one thread do not overlap`() {
withDatabase("overlap") { database ->
transaction(database) {
exec("""CREATE ALIAS IF NOT EXISTS SLEEP FOR "com.cashwu.todo.SlowSql.sleep"""")
}
onOneThread {
val blocking = measureMillis {
coroutineScope {
repeat(4) {
launch {
suspendTransaction(database) {
exec("SELECT SLEEP(200)") { rows -> rows.next() }
}
}
}
}
}
val suspending = measureMillis {
coroutineScope { repeat(4) { launch { delay(200) } } }
}
assertTrue(blocking > 700, "四個阻塞交易只花了 $blocking 毫秒")
assertTrue(suspending < 400, "四個 delay 花了 $suspending 毫秒")
}
}
}
H2 沒有內建的 SLEEP 函式,所以借用 CREATE ALIAS 掛一個自己寫的靜態方法上去,它跟 measureMillis 都在 TransactionTest.kt 的最外層
object SlowSql {
@JvmStatic
fun sleep(millis: Long): Long {
Thread.sleep(millis)
return millis
}
}
private inline fun measureMillis(block: () -> Unit): Long {
val start = System.nanoTime()
block()
return (System.nanoTime() - start) / 1_000_000
}
Thread.sleep 在這裡不是造假,它模擬的就是「JDBC driver 在等資料庫回話」那段時間,driver 等的時候 thread 一樣什麼都做不了
跑出來是 829 毫秒對 211 毫秒,同一條 thread、同樣 4 個協程、同樣每個 200 毫秒,一組是 4 乘 200,一組是 200
delay 會讓出 thread,executeQuery() 不會,suspend 這個關鍵字只代表這個函式可以被暫停,不代表它裡面每一件事都會暫停
阻塞不只是慢,配上有限的連線池還會直接互相卡住,這個測試同樣放 TransactionTest.kt
@Test
fun `one connection and one thread is enough to deadlock`() {
val source = HikariDataSource(
HikariConfig().apply {
jdbcUrl = "jdbc:h2:mem:deadlock"
driverClassName = "org.h2.Driver"
username = "sa"
password = ""
maximumPoolSize = 1
connectionTimeout = 250
poolName = "deadlock-pool"
}
)
source.use { dataSource ->
val database = dataSource.connectAndSeed()
val failure = assertFailsWith<ExposedSQLException> {
onOneThread {
launch {
suspendTransaction(database) {
Todos.selectAll().count()
delay(100)
}
}
launch {
suspendTransaction(database) {
Todos.selectAll().where { Todos.id eq 1 }.count()
}
}
}
}
assertContains(failure.message.orEmpty(), "Connection is not available")
}
}
這個測試自己組 HikariDataSource 而不是用 todoDataSource,因為要把 connectionTimeout 調到 250 毫秒,HikariCP 預設是 30 秒,配上前面說過的重試,這種測試沒人想等
第 1 個協程開交易,拿走唯一那條連線、然後 delay(100) 讓出 thread,第 2 個協程接手這條唯一的 thread,開交易,跟連線池要連線,池子是空的,於是它阻塞在 getConnection() 上,第 1 個協程 100 毫秒後想回來繼續,但唯一的 thread 被第 2 個佔著,而第 2 個在等第 1 個還連線,兩邊互相等,誰都動不了
最後是 HikariCP 的 timeout 把它拆開
Transaction attempt #0 failed: java.sql.SQLTransientConnectionException: deadlock-pool - Connection is not available, request timed out after 252ms (total=1, active=1, idle=0, waiting=0). Statement(s): SELECT TODOS.ID, TODOS.TITLE, TODOS.DONE, TODOS.CREATED_AT FROM TODOS WHERE TODOS.ID = ?
total=1, active=1, idle=0 3 個數字寫得很清楚,池子裡就一條連線、正在用、沒有閒的,後面還有 #1、#2 2 行幾乎一樣的,3 次是前面說過的上限,SQLException 會讓整個 block 重跑,而每一次都會踩到同一個死結,250 是設定值,實際等到的是 252 到 257 毫秒,把 connectionTimeout 拿掉、用 HikariCP 預設的 30 秒跑過一次,光是第 1 次嘗試就等了 30003 毫秒
這個組合看起來很極端,一條連線加一條 thread,實際上它只是把比例縮小了,連線池 5 條、event loop 12 條的專案照樣會發生,只要同時進來的請求夠多,剛好每條 thread 都在等連線的時候,整台服務就停在那裡
delay 在交易裡面看起來很蠢,但把它換成「呼叫外部 API」就一點都不蠢了,而那正是 suspendTransaction 唯一多給你的能力。這是它最需要小心的用法,交易開著的時候去等一個你控制不了的東西
上一節是設計出來的極端例子,這節換成真的跑起來的服務
Ktor Netty 拿來跑 handler 的那組 thread 由 callGroupSize 決定,ApplicationEngine.Configuration 裡的預設值就是 parallelism,也就是 availableProcessors(),所以是 12 條,連線池是 application.yaml 裡的 poolSize: 5,12 條 thread 搶 5 條連線,這個比例是這節的重點
要壓測的路由裡面只有一句 SELECT SLEEP(500),suspendTransaction 那節的 SlowSql 放在 test 底下,正式程式看不到,先複製一份到 src/main/kotlin/com/cashwu/todo/SlowSql.kt,內容跟 TransactionTest.kt 裡那個 object 一模一樣,這個檔案跟下面掛 alias 那段、/slow 路由都不是一次性的,day 22 加上 withContext 之後要用同一組東西再量一次,所以這篇量完不刪
alias 則在服務啟動時掛一次,位置在 Application.kt 的 module() 裡,day 20 在 dependencies { } 區塊後面用 by dependencies 拿了 config、repository、database 3 個東西,database 那時候只是先取出來確認容器建得起來,掛 alias 的 transaction { } 就接在它後面、install(CallId) 前面,這樣每個請求進來 SLEEP 都已經在了,不用在路由裡每次 CREATE ALIAS
fun Application.module() {
dependencies {
// day 20 的 provide 都不動
}
val config: TodoConfig by dependencies
val repository: TodoRepository by dependencies
val database: Database by dependencies
transaction(database) {
exec("""CREATE ALIAS IF NOT EXISTS SLEEP FOR "com.cashwu.todo.SlowSql.sleep"""")
}
// ...
}
transaction 要 import org.jetbrains.exposed.v1.jdbc.transactions.transaction,day 20 的 Application.kt 還沒用到它。database 上面那個 @Suppress("UNUSED_VARIABLE") 現在可以拿掉,它真的被用到了,這個變數會一直留到 day 22 清掉 alias 為止
路由跟前面的 /thread 一樣放在 routing { } 裡,差別是 /thread 測完就刪,這個要留給 day 22 當對照組
get("/slow") {
transaction(database) {
exec("SELECT SLEEP(500)") { rows -> rows.next() }
}
call.respondText("slow done")
}
./gradlew run 起 Netty,先單獨打一次確認 SLEEP 真的在睡。-w '%{time_total}' 讓 curl 把這次請求的總耗時印出來,單位是秒
curl -s -w ' %{time_total}\n' localhost:8080/slow
slow done 0.510297
再打 10 次 GET / 當基準,就是 day 02 那個回 Hello, Ktor! 的路由,沒有資料庫,這個迴圈跟下面壓測那段都是 bash 語法,直接貼進 fish 會被拒絕,貼進 zsh 多行貼上也不好控制,所以都存成檔案用 bash 跑,這段存成專案根目錄的 ping.sh,day 22 會原封不動再跑一次
for i in $(seq 1 10); do
curl -s -o /dev/null -w '%{time_total} ' localhost:8080/
sleep 0.3
done
echo
bash ping.sh
0.003165 0.003021 0.004264 0.003001 0.003977 0.002883 0.003599 0.003735 0.006173 0.004525
閒置的時候 10 次都在 7 毫秒以內,接著把 /slow 打滿,xargs -P 40 讓 40 個 curl 同時跑,總共 80 個請求,前後各記一次時間戳算出整批的 wall time,macOS 的 date 拿不到毫秒,所以借 python 拿,這段存成 slow.sh,day 22 也會用它,路徑那時候會改成參數
start=$(python3 -c 'import time; print(int(time.time() * 1000))')
seq 1 80 | xargs -P 40 -I{} curl -s -o /dev/null localhost:8080/slow
end=$(python3 -c 'import time; print(int(time.time() * 1000))')
echo "batch wall = $((end - start)) ms"
bash slow.sh
這段跑下去的同時,另開一個 shell 跑 bash ping.sh,兩邊的輸出擺在一起
batch wall = 8665 ms
2.553652 4.212209 0.003172 0.001994 0.001819 0.002355 0.003062 0.002276 0.003720 0.003062
先看第 1 行,80 個請求乘 500 毫秒,除以連線池的 5,是 8000 毫秒,實際量到 8665,多出來的是排隊跟 HTTP 本身的成本,event loop 有 12 條 thread,連線只有 5 條,任何時刻最多 5 條 thread 在等資料庫回話,其餘的在 getConnection() 上等連線,12 條大部分時間都不是在工作而是在等
第 2 行是那 10 次 GET /,前 2 次 2.5 秒跟 4.2 秒,第 3 次之後回到毫秒級,再跑兩輪,batch wall 是 8660 跟 8617 毫秒,GET / 最慢的一次分別是 3.7 秒跟 3.2 秒,每次的數字都不同,但形狀很穩定,總是開頭幾次卡住、其餘毫秒級
差別在於這個請求的連線被 Netty 分到哪一條 event loop thread,分到還沒被佔住的就毫秒級回來,分到正卡在 getConnection() 上的,就得排隊等前面那些 500 毫秒的查詢做完,同一個路由,同一份程式碼,回應時間差了 3 個數量級
day 19 講為什麼選 synchronized 的時候寫過一句「被擋住的那條 thread 在等鎖的期間什麼都不能做,包括去服務別的請求」,這裡是同一句話換一個主角,被擋住的不是鎖而是 JDBC
修法 day 22 再處理,方向就是把這段工作丟到別的 thread pool 上,讓 event loop 空出來,這一節加的東西全部留著,day 22 改完要用同一組東西再量一次才有對照
update 跟 deleteWhere 都寫完測過了,API 上卻還是只有 2 個 GET 加一個 POST,這是刻意的
要開 PUT /todos/{id} 得先讓 TodoRepository 這個介面長出 update 跟 delete 2 個方法,然後在 InMemoryTodoRepository 上實作一次,而 day 22 的題目就是把這個實作換成 Exposed 版,也就是說今天寫的那 2 個記憶體版方法,明天就會被丟掉,中間還隔著 synchronized 要不要繼續,read-modify-write 在記憶體版怎麼處理這些問題,而它們的答案在資料庫版本裡根本不存在
所以端點留給 day 22,跟 repository 一起換過去比較省事,day 12 那個 put todos responds method not allowed 也就繼續留著,那篇原話是「405 的語意本身還想看住,所以不是刪掉而是換一個還沒實作的方法,對 /todos 發 PUT 斷言 405,等哪天 PUT 也變成真端點再讓位一次」
day 20 說過 Relix 那個系列沒有走到資料庫,資料一直放在 LinkedHashMap 裡,但它在 2 個地方剛好講到這篇的重點
day 29 講 in-memory repository 的並行問題,原話是「update() 的「讀出來 → copy → 寫回去」是一段 read-modify-write,就算換成 ConcurrentHashMap 也還是有 lost update 的問題,要用 compute() 這類原子操作或直接加鎖」,接著是「真實應用要嘛用 ConcurrentHashMap + 原子操作,要嘛就交給具有交易語意的資料庫 repository」。前面那個把條件放進 WHERE 的做法,就是後半句的實際長相,而且 SQL 的 UPDATE ... WHERE 跟 ConcurrentHashMap.compute() 是同一個念頭,判斷跟修改要在同一個不可分割的動作裡完成
day 31 談效能改善的時候寫了「不能只替換 dispatcher,阻塞式資料庫 driver 仍可明確移到 Dispatchers.IO」,配的範例是 withContext(Dispatchers.IO) { userRepository.findAll() },下面還補了一句「Dispatchers.IO 讓阻塞 I/O 跑在另一組並行配額上,不佔用 Dispatchers.Default 的額度」
DSL 那部分沒問題,transaction 的用法有幾個地方要補
insert、select、update、deleteWhere 這 4 組就是正式專案在用的東西,Exposed 的 DSL 沒有「示範版」跟「正式版」的差別,要補的是這篇沒碰到的部分,batchInsert 處理大量寫入、upsert 處理「有就更新沒有就新增」、join 處理多張表,還有 Todos.selectAll() 一次撈全表這件事,資料一多就得配上分頁
transaction { } 包多大是個真問題,這篇每個測試都是一個交易包幾句 SQL,正式的寫法會是「一個 HTTP 請求對應一個交易」或者「一個 service 方法對應一個交易」,邊界要固定,不能東包一段西包一段,交易開太久會壓住連線池,開太短又保護不到該一起成功或一起失敗的東西
前面那個重試 3 次的行為,正式環境要想清楚,block 裡面只有 SQL 的話重試是好事,有寄信、送 MQ、寫檔案這種副作用就會做 2 次,要嘛把副作用移到交易外面,要嘛把 maxAttempts 設成 1 自己處理
readOnly = true 這個參數不要當成保護,H2 完全不理它,換一個資料庫行為又不一樣,真的要防止某條路徑寫入的話,做在資料庫使用者的權限上比較實在
最後是這篇量出來的那件事,現在的 todo-api 還沒有把任何一個端點接上資料庫,所以那 4.2 秒是實驗做出來的,不是線上的狀況,但 day 22 一接上去它就是真的了,這也是為什麼那篇要先處理阻塞 I/O 再談別的
Todos.insert { } 回傳的是 statement,statement[Todos.id] 拿到 AUTO_INCREMENT 發的 id、resultedValues 拿到整列,selectAll() 接 where、orderBy、limit 組出查詢,指定欄位用 select(vararg),count() 回的是 Long
update 跟 deleteWhere 回傳筆數,改不到刪不到都是 0,eq 跟 and 在 1.x 都要 import top-level 版本,eq 是因為 where 的 lambda receiver 被拿掉、成員版碰不到,and 則是本來就沒有成員版、只是換了套件名
8 條 thread 搶同一筆待辦,先讀再寫的版本 8 個都以為自己搶到,把條件搬進 UPDATE ... WHERE (ID = 2) AND (DONE = FALSE) 之後永遠是 1 個 1、7 個 0,day 19 欠的 lost update 用資料庫的交易語意結掉,transaction { } 正常回來 commit、丟 Throwable rollback,SQLException 會讓整個 block 重跑,總共最多 3 次
巢狀的 transaction { } 預設沿用外層交易,沒有獨立的 commit、rollback 或 savepoint,useNestedTransactions 是 false 時,內層丟例外被外層接住之後那筆照樣進表,打開之後才會用 savepoint 分開,隔離等級預設 READ COMMITTED,readOnly = true 只是給 driver 的提示,H2 照寫不誤
suspendTransaction 跟 transaction 的差別只在 statement 是 suspend lambda,withTransactionContext 只有在沒有 dispatcher 的時候才補 Dispatchers.IO,所以在 Netty 上它跟 transaction 一樣跑在 eventLoopGroupProxy 那條 thread
一條 thread 上 4 個 200 毫秒的交易花了 829 毫秒,4 個 delay(200) 則是 211 毫秒,數字每次跑都不一樣,重點是這 2 組差了 4 倍,一條連線加一條 thread 加一個 delay,2 個協程就能互相卡到 HikariCP timeout
40 個併發打滿之後,回 Hello, Ktor! 的路由 10 次裡有 2 次超過 2 秒,那一輪最慢的一次是 4.2 秒,閒置時則是幾毫秒的等級,PUT 跟 DELETE 端點留給 day 22,day 12 那個 405 繼續看住
下一篇把 InMemoryTodoRepository 換成 Exposed 版本,TodoRepository 這個介面該長什麼樣,transaction { } 應該在 repository 裡面還是外面,connection pool 要開多大才對得上 Dispatchers.IO 的並行度,還有這篇那 4.2 秒接上 Dispatchers.IO 之後會變成多少,DAO 那條路跟這篇用的 DSL 差在哪裡
同步刊登於 Blog
圖片來源:AI 產生