GRDB.swift 与 Combine 集成指南:用响应式发布器进行数据库异步读写与实时观察
【免费下载链接】GRDB.swiftA toolkit for SQLite databases, with a focus on application development项目地址: https://gitcode.com/GitHub_Trending/gr/GRDB.swift
本指南面向在 Apple 平台(macOS / iOS / tvOS / watchOS)上使用 Combine 框架的开发者,系统讲解 GRDB.swift 提供的一系列数据库发布器(Database Publishers):如何用readPublisher/writePublisher异步执行数据库读写、用migratePublisher驱动数据库迁移、用ValueObservation.publisher/DatabaseRegionObservation.publisher观察数据库值与事务变更,以及如何避免在组合多个发布器时破坏数据一致性。读完本文,你将掌握一套"发布—订阅"式的 SQLite 数据访问与变更通知方案,并理解其底层调度与事务语义。
GRDB.swift 在支持 Combine 框架的系统上,提供了把数据库值与事件发布为 Combine 发布器的能力。所有相关发布器都被组织在DatabasePublishers命名空间中(见 GRDB/Core/DatabasePublishers.swift),并通过DatabaseReader、DatabaseWriter、ValueObservation、DatabaseRegionObservation等类型对外暴露构造入口。
快速上手:五种典型用法
在使用下面的发布器之前,请先参照 数据库连接 文档完成DatabaseQueue、DatabasePool或DatabaseSnapshot的初始化。
异步读取数据库
readPublisher从数据库读取一个值并将其发布出去,发布成功后即完成:
// DatabasePublishers.Read<[Player]> let players = dbQueue.readPublisher { db in try Player.fetchAll(db) }异步写入数据库
writePublisher在数据库事务内执行更新,并发布一个结果值:
// DatabasePublishers.Write<Void> let write = dbQueue.writePublisher { db in try Player(...).insert(db) } // DatabasePublishers.Write<Int> let newPlayerCount = dbQueue.writePublisher { db -> Int in try Player(...).insert(db) return try Player.fetchCount(db) }异步迁移数据库
DatabaseMigrator也提供了发布器版本,用于在订阅时执行一次异步迁移:
// DatabasePublishers.Migrate let migrator: DatabaseMigrator = ... let publisher = migrator.migratePublisher(dbQueue)观察数据库值的变化
ValueObservation.tracking配合.publisher(in:),会在数据库发生变化时持续发布最新值:
// 一个输出 [Player]、失败类型为 Error 的发布器 let publisher = ValueObservation .tracking { db in try Player.fetchAll(db) } .publisher(in: dbQueue) // 一个输出 Int?、失败类型为 Error 的发布器 let publisher = ValueObservation .tracking { db in try Int.fetchOne(db, sql: "SELECT MAX(score) FROM player") } .publisher(in: dbQueue)观察影响数据库区域的事务
DatabaseRegionObservation的发布器在"被跟踪区域"被事务影响时,发布一个数据库连接对象:
// 一个输出 Database、失败类型为 Error 的发布器 let publisher = DatabaseRegionObservation .tracking(Player.all()) .publisher(in: dbQueue) let cancellable = publisher.sink( receiveCompletion: { completion in ... }, receiveValue: { (db: Database) in print("Exclusive write access to the database after players have been impacted") }) // 一个输出 Database、失败类型为 Error 的发布器 let publisher = DatabaseRegionObservation .tracking(SQLRequest<Int>(sql: "SELECT MAX(score) FROM player")) .publisher(in: dbQueue) let cancellable = publisher.sink( receiveCompletion: { completion in ... }, receiveValue: { (db: Database) in print("Exclusive write access to the database after maximum score has been impacted") })异步数据库访问(Asynchronous Database Access)
GRDB 提供以下执行异步数据库访问的发布器:
readPublisher(receiveOn:value:)writePublisher(receiveOn:updates:)writePublisher(receiveOn:updates:thenRead:)migratePublisher(_:receiveOn:)
这类发布器遵循共同的契约:数据库访问不会在订阅之前发生,每次订阅都会启动一次全新的数据库访问;发布器可以在任意线程被订阅,并默认在主队列(main queue)上发布值与完成事件,除非通过receiveOn参数指定其他 Combine Scheduler。
从源码实现看(GRDB/Core/DatabaseReader.swift、GRDB/Core/DatabaseWriter.swift),这些发布器内部都基于OnDemandFuture(见 GRDB/Utils/OnDemandFuture.swift)构造"每次订阅才执行"的 Future,再通过receiveValues(on:)切回指定调度器,最终包装为类型化的DatabasePublishers.Read或DatabasePublishers.Write。注释中也特别强调:receiveValues(on:)的目的就是避免用户在数据库派发队列上处理发布值。
DatabaseReader.readPublisher(receiveOn:value:)
返回一个在异步获取数据库值之后完成的发布器:
// DatabasePublishers.Read<[Player]> let players = dbQueue.readPublisher { db in try Player.fetchAll(db) }要点:
- 只读约束:任何修改数据库的尝试都会让订阅以错误结束。结合源码看,
readPublisher内部调用asyncRead(GRDB/Core/DatabaseReader.swift),底层走DatabaseReader的异步只读通道;文档注释明确,对连接的任何写入都会抛出 resultCode 为SQLITE_READONLY的DatabaseError。 - 队列与快照:使用数据库队列或数据库快照时,本次读取需要等待该队列/快照上可能存在的其他并发数据库访问结束。
- 数据库池:使用数据库池时,读取通常是非阻塞的,除非已达到最大并发读数量——此时读操作必须等待其他读操作完成。该上限可通过 DatabasePool 配置调整。
- 连接有效性:闭包中的
Database参数仅在闭包执行期间有效,不能保存或返回供以后使用(源码 GRDB/Core/DatabaseReader.swift 的文档注释明确警告了这一点)。 - 事务隔离:读取操作在事务中被隔离,看不到其他进程或并发写操作带来的未提交变更。
- 调度:默认在主队列发布,可通过
receiveOn传入自定义 Scheduler;默认值在源码中即为DispatchQueue.main(GRDB/Core/DatabaseReader.swift)。
DatabaseWriter.writePublisher(receiveOn:updates:)
返回一个在数据库更新于事务内成功执行后完成的发布器:
// DatabasePublishers.Write<Void> let write = dbQueue.writePublisher { db in try Player(...).insert(db) } // DatabasePublishers.Write<Int> let newPlayerCount = dbQueue.writePublisher { db -> Int in try Player(...).insert(db) return try Player.fetchCount(db) }要点:
- 事务语义:闭包中的全部操作被包裹在一个事务中;如果闭包抛出错误,事务回滚,发布器以该错误完成(源码注释见 GRDB/Core/DatabaseWriter.swift)。并发数据库访问永远看不到部分更新的中间状态,即使更新来自其他进程也一样。
- 串行化:写操作被异步派发到 writer 的写派发队列,与该
DatabaseWriter执行的所有数据库更新串行化。 - 调度:默认在主队列完成,可用
receiveOn指定其他 Scheduler;可被任意线程订阅,每次订阅启动一次新的数据库访问。
当使用数据库池,且应用执行"一批数据库更新 + 一些较慢的后续读取"时,可以考虑下文的writePublisher(receiveOn:updates:thenRead:)获得调度优化。
DatabaseWriter.writePublisher(receiveOn:updates:thenRead:)
返回一个在事务内更新成功执行、随后又完成取值后才完成的发布器:
// DatabasePublishers.Write<Int> let newPlayerCount = dbQueue.writePublisher( updates: { db in try Player(...).insert(db) } thenRead: { db, _ in try Player.fetchCount(db) })它发布的值与writePublisher(receiveOn:updates:)完全一致:
// DatabasePublishers.Write<Int> let newPlayerCount = dbQueue.writePublisher { db -> Int in try Player(...).insert(db) return try Player.fetchCount(db) }两者的区别在于:后一批"读取"操作被移入thenRead闭包执行。thenRead接受两个参数——一个只读数据库连接,以及updates闭包的返回值;这允许你在两个闭包之间传递信息(上述示例中该参数被忽略)。源码 GRDB/Core/DatabaseWriter.swift 展示了其实现细节:updates在db.inTransaction中执行并提交,随后通过spawnConcurrentRead派发一个并发读,让thenRead读取到更新后的数据库状态。
要点:
- 数据库池下的调度优化:使用数据库池时,
thenRead能看到updates留下的数据库状态,但不会阻塞任何并发写,从而降低写竞争(write contention)。 - 数据库队列下的行为:使用数据库队列时结果保证一致,但不应用上述调度优化。
- 与其它写发布器一样:可任意线程订阅、每次订阅启动一次新的数据库访问、默认在主队列完成(可用
receiveOn调整)。
DatabaseMigrator.migratePublisher(_:receiveOn:)
返回一个异步迁移数据库的发布器:
let migrator: DatabaseMigrator = ... let publisher = migrator.migratePublisher(dbQueue)其源码实现(GRDB/Migration/DatabaseMigrator.swift)在订阅时调用asyncMigrate(writer),成功后发布单个事件并完成;值与完成默认在主队列发布,可通过receiveOn指定其他 Scheduler。这在启动流程中把"打开数据库并迁移"以发布器形态接入 Combine 管线(例如与其它异步初始化任务组合)非常方便。
数据库观察(Database Observation)
数据库观察发布器基于 ValueObservation 与 DatabaseRegionObservation 构建。如果你的应用需要的是非 Combine 形态的变更通知,请参阅 数据库变更观察 相关章节。
相关入口:
ValueObservation.publisher(in:scheduling:)SharedValueObservation.publisher()DatabaseRegionObservation.publisher(in:)
ValueObservation.publisher(in:scheduling:)
ValueObservation用于跟踪数据库值的变化,可以很方便地转为 Combine 发布器:
let observation = ValueObservation.tracking { db in try Player.fetchAll(db) } // 一个输出 [Player]、失败类型为 Error 的发布器 let publisher = observation.publisher(in: dbQueue)该发布器的行为与ValueObservation一致(源码实现见 GRDB/ValueObservation/ValueObservation.swift,发布器类型为DatabasePublishers.Value):
- 先发初值:在后续变化之前,会先通知一个初始值。
- 合并变化:可能会把随后发生的多次变化合并为一次通知。
- 可能发布连续相同值:可以用 Combine 的
removeDuplicates()操作符过滤多余重复值,但更推荐使用 GRDB 自带的 removeDuplicates() 操作符。 - 只随取消而完成:发布器仅在订阅被取消时结束。
- 默认异步主线程通知:默认在主线程异步地通知初始值、后续变化与错误。
这个默认调度行为可以通过scheduling参数配置。注意:scheduling不接受 Combine Scheduler,而接受 GRDB 的 ValueObservationScheduler。
例如,.immediate调度器保证初始值在订阅发生时立即被通知,从而让 UI 不必等待任何异步通知即可更新:
// 初始值的即时通知 let cancellable = observation .publisher( in: dbQueue, scheduling: .immediate) // <- .sink( receiveCompletion: { completion in ... }, receiveValue: { (players: [Player]) in print("Fresh players: \(players)") }) // <- 在这里 "fresh players" 已经被打印了。注意:.immediate调度器要求发布器从主线程订阅,否则会触发致命错误(fatal error)。这一点在源码文档注释中亦有强调(GRDB/ValueObservation/ValueObservation.swift),源码默认值则是.async(onQueue: .main)。
SharedValueObservation.publisher()
SharedValueObservation同样跟踪数据库值的变化,并可在多个订阅者之间共享同一次数据库观察:
let sharedObservation = ValueObservation .tracking { db in try Player.fetchAll(db) } .shared(in: dbQueue) // 一个输出 [Player]、失败类型为 Error 的发布器 let publisher = sharedObservation.publisher()该发布器的行为与SharedValueObservation一致(构造入口见 GRDB/ValueObservation/SharedValueObservation.swift)。当多个 UI 组件需要观察同一份数据、但又不希望各自开启独立的观察开销时,这是更合适的选择。
DatabaseRegionObservation.publisher(in:)
DatabaseRegionObservation会通知所有影响被跟踪数据库区域的事务,可转为 Combine 发布器:
let request = Player.all() let observation = DatabaseRegionObservation.tracking(request) // 一个输出 Database、失败类型为 Error 的发布器 let publisher = observation.publisher(in: dbQueue)该发布器可在任意线程创建和订阅,并在"受保护的派发队列"(protected dispatch queue)中、与所有数据库更新串行化地递送数据库连接;只有在发生数据库错误时才会完成。下面是一个完整可运行的示例:
let request = Player.all() let cancellable = DatabaseRegionObservation .tracking(request) .publisher(in: dbQueue) .sink( receiveCompletion: { completion in ... }, receiveValue: { (db: Database) in print("Players have changed.") }) try dbQueue.write { db in try Player(name: "Arthur").insert(db) try Player(name: "Barbara").insert(db) } // 打印 "Players have changed." try dbQueue.write { db in try Player.deleteAll(db) } // 打印 "Players have changed."实现细节(GRDB/Core/DatabaseRegionObservation.swift)表明:DatabaseRegionObserver实现了TransactionObserver协议,通过observes(eventsOfKind:)与databaseDidChange(with:)判断事件是否命中被跟踪区域,并在databaseDidCommit(_:)时回调onChange(db),从而让发布器递送一个仅在被发布瞬间有效的Database连接。也正因如此,源码文档明确建议:不要用receive(on:options:)或任何会调度发布元素的Publisher方法去重新调度该发布器(GRDB/Core/DatabaseRegionObservation.swift),也不要保存或返回该连接供以后使用。更多信息参见 DatabaseRegionObservation。
Combine 与数据一致性(Combine and Data Consistency)
当你用combineLatest、zip等 Combine 操作符组合多个数据库发布器时,会失去所有数据一致性保证。
原因在于:每个数据库发布器彼此隔离,各自看到数据库的某一状态。一旦数据库变化发生在各发布器操作之间,各发布器处理或发布的值就可能互相不匹配。
换言之:只要某项数据库访问或观察依赖于某个数据库不变式(database invariant),就应该只定义一个数据库发布器,而不是组合多个发布器。这样才能力求避免并发写打乱应用状态、引入难以排查的 bug。好消息是:所有数据库发布器都支持在一个闭包内执行多次请求。
下面的例子中,HallOfFame值由唯一一个发布器产生,因此可以保证发布出来的值永远不会自相矛盾:
struct HallOfFame { // 不变式:bestPlayers.count <= totalPlayerCount var totalPlayerCount: Int var bestPlayers: [Player] } // 正确:数据一致性有保证 let hallOfFamePublisher = ValueObservation .tracking { db -> HallOfFame in // 第 1 个请求 let totalPlayerCount = try Player.fetchCount(db) // 第 2 个请求 let bestPlayers = try Player .order(\.score.desc) .limit(10) .fetchAll(db) // 100% 有保证 assert(bestPlayers.count <= totalPlayerCount) // 把结果合并到一起 return HallOfFame( totalPlayerCount: totalPlayerCount, bestPlayers: bestPlayers) } .publisher(in: dbQueue)对比下面这个把两个数据库发布器组合在一起的错误版本:
// 单独看没问题 let totalPlayerCountPublisher = ValueObservation .tracking(Player.fetchCount) .publisher(in: dbQueue) // 单独看没问题 let bestPlayerPublisher = ValueObservation .tracking(Player .order(\.score.desc) .limit(10) .fetchAll) .publisher(in: dbQueue) // 错误:数据一致性没有保证 let hallOfFamePublisher = totalPlayerCountPublisher .combineLatest(bestPlayerPublisher) .map(HallOfFame.init(totalPlayerCount:bestPlayers)) let cancellable = hallOfFamePublisher.sink( receiveCompletion: { completion in ... }, receiveValue: { hallOfFame in // 若某些玩家恰好在错误的时间点被删除, // 这个断言可能失败 assert(hallOfFame.bestPlayers.count <= hallOfFame.totalPlayerCount) })发布器的类型与实现结构
所有数据库发布器都声明在DatabasePublishers命名空间下(见 GRDB/Core/DatabasePublishers.swift),按用途分为以下几类,各自遵循Publisher协议且Failure == Error:
| 发布器类型 | 输出(Output) | 构造入口 | 行为 |
|---|---|---|---|
DatabasePublishers.Read<Output> | 单个值 | DatabaseReader.readPublisher | 发布一个读取值后完成 |
DatabasePublishers.Write<Output> | 单个值 | DatabaseWriter.writePublisher(两个重载) | 事务内更新、可选后读,然后完成 |
DatabasePublishers.Migrate | 单个事件 | DatabaseMigrator.migratePublisher | 执行迁移后完成 |
DatabasePublishers.Value<Output> | 持续的值序列 | ValueObservation.publisher、SharedValueObservation.publisher | 初始值 + 持续变化 |
DatabasePublishers.DatabaseRegion | Database | DatabaseRegionObservation.publisher | 每次受影响事务提交时递送一次连接 |
从类型定义看(如 GRDB/Core/DatabaseReader.swift 中的Read、GRDB/Core/DatabaseWriter.swift 中的Write),这些公开类型只是对内部上游AnyPublisher的类型擦除包装,真正的数据库访问逻辑(asyncRead、asyncWrite、asyncMigrate、start(in:scheduling:onError:onChange:))都被封装在闭包中、延迟到订阅时刻执行。
测试与验证
仓库的 Tests/GRDBTests/GRDBCombineTests 目录为上述发布器提供了完整的行为测试,是理解 API 契约的最好参考:
- DatabaseReaderReadPublisherTests.swift:覆盖
readPublisher的取值、错误传播与调度行为; - DatabaseWriterWritePublisherTests.swift:覆盖
writePublisher两个重载的事务、thenRead传递与调度优化; - ValueObservationPublisherTests.swift:覆盖
ValueObservation.publisher的初始值、合并通知、.immediate调度等; - DatabaseRegionObservationPublisherTests.swift:覆盖
DatabaseRegionObservation.publisher的事务通知行为。
阅读这些测试可以更直观地看到各发布器在"正常提交""回滚""错误抛出""订阅取消"等场景下的预期行为,便于你在自己的工程中写出符合 GRDB 语义的 Combine 代码。
小结
- 一次性异步访问:
readPublisher提供只读、事务隔离的异步读取;writePublisher提供事务包裹的异步写入,其thenRead变体在数据库池上具备降低写竞争的调度优化;migratePublisher把数据库迁移接入 Combine 管线。 - 持续观察:
ValueObservation.publisher(in:scheduling:)与SharedValueObservation.publisher()持续发布数据库值,DatabaseRegionObservation.publisher(in:)在被跟踪区域受影响时递送Database连接;默认都在主队列异步通知,.immediate调度可用于 UI 的即时初值呈现。 - 一致性铁律:涉及多个请求且依赖不变式的场景,务必收敛为一个发布器内完成全部请求,避免用
combineLatest/zip组合多个数据库发布器而破坏一致性。
将这套发布器接入你的 Combine 数据流后,SQLite 的读写与变更通知便与 SwiftUI、UIKit 的响应式界面自然衔接,同时保持 GRDB 一贯的强类型、事务安全与性能特征。
【免费下载链接】GRDB.swiftA toolkit for SQLite databases, with a focus on application development项目地址: https://gitcode.com/GitHub_Trending/gr/GRDB.swift
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考