iT邦幫忙

2026 iThome 鐵人賽

DAY 21
0
Software Development

Kotlin Ktor 實戰 101系列 第 21 篇

Kotlin Ktor 實戰 101 Day 21 Exposed CRUD 與 Transaction

  • 分享至 

  • xImage
  •  

https://ithelp.ithome.com.tw/upload/images/20260909/20121948LfQo3aVCtG.jpg

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,跟它們送出去的 SQL
  • update 回傳的筆數就是併發的答案,day 19 欠的那個 lost update
  • transaction { } 的邊界,commit 跟 rollback 各發生在哪一行
  • 巢狀 transaction 預設沿用外層交易,沒有獨立的 commit 與 rollback
  • isolation level 跟 readOnly 是連線上的 2 個參數,其中一個 H2 不理你
  • suspendTransaction 掛著 suspend,但沒有換掉任何一條 thread
  • 一條連線加一條 thread,2 個協程就能互相卡死
  • 打滿之後,連 Hello, Ktor! 都要等 3 秒半
  • PUT 跟 DELETE 這篇不補,day 12 那個 405 再撐一天
  • 跟 Relix 的對照,那個系列停在哪一句
  • 這樣的 CRUD 跟 transaction 能不能上線

四組 DSL 跟它們送出去的 SQL

這篇新增 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 拿掉了

update 回傳的筆數就是併發的答案

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 { } 的邊界在哪裡

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 { } 裡面再包一個 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 在高併發下是有成本的

isolation level 跟 readOnly

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

suspendTransaction 沒有讓任何東西變成非阻塞

前面每一個 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 這個關鍵字只代表這個函式可以被暫停,不代表它裡面每一件事都會暫停

一條連線加一條 thread 就能卡死

阻塞不只是慢,配上有限的連線池還會直接互相卡住,這個測試同樣放 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 唯一多給你的能力。這是它最需要小心的用法,交易開著的時候去等一個你控制不了的東西

打滿之後連 Hello, Ktor! 都要等

上一節是設計出來的極端例子,這節換成真的跑起來的服務

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 改完要用同一組東西再量一次才有對照

PUT 跟 DELETE 這篇不補

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 也變成真端點再讓位一次」

跟 Relix 的對照

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 的額度」

這樣的 CRUD 跟 transaction 能不能上線

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 產生


上一篇
Kotlin Ktor 實戰 101 Day 20 Exposed 入門與資料庫連線
系列文
Kotlin Ktor 實戰 101 共 21 篇
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言