diff --git a/spring-tx/src/main/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensions.kt b/spring-tx/src/main/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensions.kt index 3a874b860a..b819fee09a 100644 --- a/spring-tx/src/main/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensions.kt +++ b/spring-tx/src/main/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensions.kt @@ -7,23 +7,24 @@ import kotlinx.coroutines.reactive.asFlow import kotlinx.coroutines.reactive.awaitFirstOrNull import kotlinx.coroutines.reactor.asFlux import kotlinx.coroutines.reactor.mono +import org.springframework.transaction.ReactiveTransaction /** - * Coroutines variant of [TransactionalOperator.transactional] with a [Flow] parameter. + * Coroutines variant of [TransactionalOperator.transactional] as a [Flow] extension. * * @author Sebastien Deleuze * @since 5.2 */ @ExperimentalCoroutinesApi -fun TransactionalOperator.transactional(flow: Flow): Flow = - transactional(flow.asFlux()).asFlow() +fun Flow.transactional(operator: TransactionalOperator): Flow = + operator.transactional(asFlux()).asFlow() /** -* Coroutines variant of [TransactionalOperator.transactional] with a suspending lambda +* Coroutines variant of [TransactionalOperator.execute] with a suspending lambda * parameter. * * @author Sebastien Deleuze * @since 5.2 */ -suspend fun TransactionalOperator.transactional(f: suspend () -> T?): T? = - transactional(mono(Dispatchers.Unconfined) { f() }).awaitFirstOrNull() +suspend fun TransactionalOperator.executeAndAwait(f: suspend (ReactiveTransaction) -> T?): T? = + execute { status -> mono(Dispatchers.Unconfined) { f(status) } }.awaitFirstOrNull() diff --git a/spring-tx/src/test/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensionsTests.kt b/spring-tx/src/test/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensionsTests.kt index ecbe883cb8..be93faf06d 100644 --- a/spring-tx/src/test/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensionsTests.kt +++ b/spring-tx/src/test/kotlin/org/springframework/transaction/reactive/TransactionalOperatorExtensionsTests.kt @@ -34,7 +34,7 @@ class TransactionalOperatorExtensionsTests { fun commitWithSuspendingFunction() { val operator = TransactionalOperator.create(tm, DefaultTransactionDefinition()) runBlocking { - operator.transactional { + operator.executeAndAwait { delay(1) true } @@ -48,7 +48,7 @@ class TransactionalOperatorExtensionsTests { val operator = TransactionalOperator.create(tm, DefaultTransactionDefinition()) runBlocking { try { - operator.transactional { + operator.executeAndAwait { delay(1) throw IllegalStateException() } @@ -72,7 +72,7 @@ class TransactionalOperatorExtensionsTests { emit(4) } runBlocking { - val list = operator.transactional(flow).toList() + val list = flow.transactional(operator).toList() assertThat(list).hasSize(4) } assertThat(tm.commit).isTrue() @@ -89,7 +89,7 @@ class TransactionalOperatorExtensionsTests { } runBlocking { try { - operator.transactional(flow).toList() + flow.transactional(operator).toList() } catch (ex: IllegalStateException) { assertThat(tm.commit).isFalse() assertThat(tm.rollback).isTrue() diff --git a/src/docs/asciidoc/languages/kotlin.adoc b/src/docs/asciidoc/languages/kotlin.adoc index 7cd8265d1b..4670eaa593 100644 --- a/src/docs/asciidoc/languages/kotlin.adoc +++ b/src/docs/asciidoc/languages/kotlin.adoc @@ -577,8 +577,51 @@ class UserHandler(builder: WebClient.Builder) { === Transactions Transactions on Coroutines are supported via the programmatic variant of the Reactive -transaction management provided as of Spring Framework 5.2. `TransactionalOperator.transactional` -extensions with suspending lambda and Kotlin `Flow` parameter are provided for that purpose. +transaction management provided as of Spring Framework 5.2. + +For suspending functions, a `TransactionalOperator.executeAndAwait` extension is provided. + +[source,kotlin,indent=0] +---- + import org.springframework.transaction.reactive.executeAndAwait + + class PersonRepository(private val operator: TransactionalOperator) { + + suspend fun initDatabase() = operator.executeAndAwait { + insertPerson1() + insertPerson2() + } + + private suspend fun insertPerson1() { + // INSERT SQL statement + } + + private suspend fun insertPerson2() { + // INSERT SQL statement + } + } +---- + +For Kotlin `Flow`, a `Flow.transactional` extension is provided. + +[source,kotlin,indent=0] +---- + import org.springframework.transaction.reactive.transactional + + class PersonRepository(private val operator: TransactionalOperator) { + + fun updatePeople() = findPeople().map(::updatePerson).transactional(operator) + + private fun findPeople(): Flow { + // SELECT SQL statement + } + + private suspend fun updatePerson(person: Person): Person { + // UPDATE SQL statement + } + } +---- + [[kotlin-spring-projects-in-kotlin]] == Spring Projects in Kotlin