Propagate the context in Coroutines transactions
This commit ensures that CoroutineContext is properly propagated in transactional suspending functions. Both annotation and functional variants are supported. Closes gh-27308
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2021 the original author or authors.
|
||||
* Copyright 2002-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -22,6 +22,8 @@ import java.util.concurrent.ConcurrentMap;
|
||||
|
||||
import io.vavr.control.Try;
|
||||
import kotlin.coroutines.Continuation;
|
||||
import kotlin.coroutines.CoroutineContext;
|
||||
import kotlinx.coroutines.Job;
|
||||
import kotlinx.coroutines.reactive.AwaitKt;
|
||||
import kotlinx.coroutines.reactive.ReactiveFlowKt;
|
||||
import org.apache.commons.logging.Log;
|
||||
@@ -363,7 +365,7 @@ public abstract class TransactionAspectSupport implements BeanFactoryAware, Init
|
||||
|
||||
InvocationCallback callback = invocation;
|
||||
if (corInv != null) {
|
||||
callback = () -> CoroutinesUtils.invokeSuspendingFunction(method, corInv.getTarget(), corInv.getArguments());
|
||||
callback = () -> KotlinDelegate.invokeSuspendingFunction(method, corInv);
|
||||
}
|
||||
Object result = txSupport.invokeWithinTransaction(method, targetClass, callback, txAttr, (ReactiveTransactionManager) tm);
|
||||
if (corInv != null) {
|
||||
@@ -883,6 +885,12 @@ public abstract class TransactionAspectSupport implements BeanFactoryAware, Init
|
||||
private static Object awaitSingleOrNull(Publisher<?> publisher, Object continuation) {
|
||||
return AwaitKt.awaitSingleOrNull(publisher, (Continuation<Object>) continuation);
|
||||
}
|
||||
|
||||
public static Publisher<?> invokeSuspendingFunction(Method method, CoroutinesInvocationCallback callback) {
|
||||
CoroutineContext coroutineContext = ((Continuation<?>) callback.getContinuation()).getContext().minusKey(Job.Key);
|
||||
return CoroutinesUtils.invokeSuspendingFunction(coroutineContext, method, callback.getTarget(), callback.getArguments());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -16,14 +16,17 @@
|
||||
|
||||
package org.springframework.transaction.reactive
|
||||
|
||||
import java.util.Optional
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.currentCoroutineContext
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.reactive.asFlow
|
||||
import kotlinx.coroutines.reactive.awaitLast
|
||||
import kotlinx.coroutines.reactor.asFlux
|
||||
import kotlinx.coroutines.reactor.mono
|
||||
import org.springframework.transaction.ReactiveTransaction
|
||||
import java.util.*
|
||||
import kotlin.coroutines.CoroutineContext
|
||||
import kotlin.coroutines.EmptyCoroutineContext
|
||||
|
||||
/**
|
||||
* Coroutines variant of [TransactionalOperator.transactional] as a [Flow] extension.
|
||||
@@ -31,8 +34,8 @@ import org.springframework.transaction.ReactiveTransaction
|
||||
* @author Sebastien Deleuze
|
||||
* @since 5.2
|
||||
*/
|
||||
fun <T : Any> Flow<T>.transactional(operator: TransactionalOperator): Flow<T> =
|
||||
operator.transactional(asFlux()).asFlow()
|
||||
fun <T : Any> Flow<T>.transactional(operator: TransactionalOperator, context: CoroutineContext = EmptyCoroutineContext): Flow<T> =
|
||||
operator.transactional(asFlux(context)).asFlow()
|
||||
|
||||
/**
|
||||
* Coroutines variant of [TransactionalOperator.execute] with a suspending lambda
|
||||
@@ -42,6 +45,8 @@ fun <T : Any> Flow<T>.transactional(operator: TransactionalOperator): Flow<T> =
|
||||
* @author Mark Paluch
|
||||
* @since 5.2
|
||||
*/
|
||||
suspend fun <T> TransactionalOperator.executeAndAwait(f: suspend (ReactiveTransaction) -> T): T =
|
||||
execute { status -> mono(Dispatchers.Unconfined) { f(status) } }.map { value -> Optional.ofNullable(value) }
|
||||
suspend fun <T> TransactionalOperator.executeAndAwait(f: suspend (ReactiveTransaction) -> T): T {
|
||||
val context = currentCoroutineContext().minusKey(Job.Key)
|
||||
return execute { status -> mono(context) { f(status) } }.map { value -> Optional.ofNullable(value) }
|
||||
.defaultIfEmpty(Optional.empty()).awaitLast().orElse(null)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user