使用协程时Micronaut Data R2DBC事务传播异常问题
Micronaut Data R2DBC 结合 Kotlin 协程时事务上下文传播异常
问题现象
使用Micronaut Data R2DBC搭配Kotlin协程时,事务上下文无法正确传播,触发NoTransactionException错误,提示**"预期存在事务,但在Reactive上下文中未找到"**。
相关代码
Repository 定义
@Transactional(Transactional.TxType.MANDATORY) @R2dbcRepository(dialect = Dialect.POSTGRES) interface RecordTransactionalCoroutineRepository : CoroutineCrudRepository<Record, UUID>
服务层代码
@Transactional open fun saveAllUsingCoroutines(records: Iterable<Record>): Flow<Record> = coroutineRepository.saveAll(records)
堆栈跟踪信息
14:37:54.217 [reactor-tcp-epoll-2] WARN i.m.d.r.o.DefaultR2dbcRepositoryOperations - 回滚事务:RecordTransactionalService.saveAllUsingCoroutines 执行出错:预期存在事务,但在Reactive上下文中未找到。数据源为default io.micronaut.transaction.exceptions.NoTransactionException: 预期存在事务,但在Reactive上下文中未找到。 at io.micronaut.data.r2dbc.operations.DefaultR2dbcRepositoryOperations.lambda$withTransaction$15(DefaultR2dbcRepositoryOperations.java:410) at reactor.core.publisher.FluxDeferContextual.subscribe(FluxDeferContextual.java:49) at reactor.core.publisher.Flux.subscribe(Flux.java:8660) at kotlinx.coroutines.reactive.PublisherAsFlow.collectImpl(ReactiveFlow.kt:94) at kotlinx.coroutines.reactive.PublisherAsFlow.collect(ReactiveFlow.kt:79) at kotlinx.coroutines.reactive.FlowSubscription.consumeFlow(ReactiveFlow.kt:275) at kotlinx.coroutines.reactive.FlowSubscription.flowProcessing(ReactiveFlow.kt:209) at kotlinx.coroutines.reactive.FlowSubscription.access$flowProcessing(ReactiveFlow.kt:187) at kotlinx.coroutines.reactive.FlowSubscription$createInitialContinuation$1$1.invoke(ReactiveFlow.kt:204) at kotlinx.coroutines.reactive.FlowSubscription$createInitialContinuation$1$1.invoke(ReactiveFlow.kt:204) at kotlin.coroutines.intrinsics.IntrinsicsKt__IntrinsicsJvmKt$createCoroutineUnintercepted$$inlined$createCoroutineFromSuspendFunction$IntrinsicsKt__IntrinsicsJvmKt$2.invokeSuspend(IntrinsicsJvm.kt:205) at kotlin.coroutines.jvm.internal.BaseContinuationImpl.resumeWith(ContinuationImpl.kt:33) at kotlinx.coroutines.internal.DispatchedContinuationKt.resumeCancellableWith(DispatchedContinuation.kt:367) at kotlinx.coroutines.internal.DispatchedContinuationKt.resumeCancellableWith$default(DispatchedContinuation.kt:278) at kotlinx.coroutines.intrinsics.CancellableKt.startCoroutineCancellable(Cancellable.kt:18) at kotlinx.coroutines.reactive.FlowSubscription$createInitialContinuation$$inlined$Continuation$1.resumeWith(Continuation.kt:162) at kotlinx.coroutines.reactive.FlowSubscription.request(ReactiveFlow.kt:267) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.request(FluxContextWrite.java:136) at reactor.core.publisher.FluxUsingWhen$UsingWhenSubscriber.request(FluxUsingWhen.java:319) at reactor.core.publisher.Operators$DeferredSubscription.set(Operators.java:1717) at reactor.core.publisher.FluxUsingWhen$UsingWhenSubscriber.onSubscribe(FluxUsingWhen.java:409) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onSubscribe(FluxContextWrite.java:101) at kotlinx.coroutines.reactive.FlowAsPublisher.subscribe(ReactiveFlow.kt:182) at reactor.core.publisher.FluxSource.subscribe(FluxSource.java:67) at reactor.core.publisher.Flux.subscribe(Flux.java:8660) at reactor.core.publisher.FluxUsingWhen$ResourceSubscriber.onNext(FluxUsingWhen.java:195) at reactor.core.publisher.Operators$BaseFluxToMonoOperator.completePossiblyEmpty(Operators.java:2034) at reactor.core.publisher.MonoHasElements$HasElementsSubscriber.onComplete(MonoHasElements.java:93) at reactor.core.publisher.MonoIgnoreElements$IgnoreElementsSubscriber.onComplete(MonoIgnoreElements.java:89) at io.r2dbc.postgresql.util.FluxDiscardOnCancel$FluxDiscardOnCancelSubscriber.onComplete(FluxDiscardOnCancel.java:104) at reactor.core.publisher.FluxPeek$PeekSubscriber.onComplete(FluxPeek.java:260) at reactor.core.publisher.FluxPeek$PeekSubscriber.onComplete(FluxPeek.java:260) at reactor.core.publisher.FluxHandle$HandleSubscriber.onComplete(FluxHandle.java:222) at io.r2dbc.postgresql.util.FluxDiscardOnCancel$FluxDiscardOnCancelSubscriber.onComplete(FluxDiscardOnCancel.java:104) at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onComplete(FluxContextWrite.java:126) at reactor.core.publisher.FluxCreate$BaseSink.complete(FluxCreate.java:460) at reactor.core.publisher.FluxCreate$BufferAsyncSink.drain(FluxCreate.java:805) at reactor.core.publisher.FluxCreate$BufferAsyncSink.complete(FluxCreate.java:753) at reactor.core.publisher.FluxCreate$SerializedFluxSink.drainLoop(FluxCreate.java:247) at reactor.core.publisher.FluxCreate$SerializedFluxSink.drain(FluxCreate.java:213) at reactor.core.publisher.FluxCreate$SerializedFluxSink.complete(FluxCreate.java:204) at io.r2dbc.postgresql.client.ReactorNettyClient$Conversation.complete(ReactorNettyClient.java:671) at io.r2dbc.postgresql.client.ReactorNettyClient$BackendMessageSubscriber.emit(ReactorNettyClient.java:937) at io.r2dbc.postgresql.client.ReactorNettyClient$BackendMessageSubscriber.onNext(ReactorNettyClient.java:813) at io.r2dbc.postgresql.client.ReactorNettyClient$BackendMessageSubscriber.onNext(ReactorNettyClient.java:719) at reactor.core.publisher.FluxHandle$HandleSubscriber.onNext(FluxHandle.java:128) at reactor.core.publisher.FluxPeekFuseable$PeekConditionalSubscriber.onNext(FluxPeekFuseable.java:854) at reactor.core.publisher.FluxMap$MapConditionalSubscriber.onNext(FluxMap.java:224) at reactor.core.publisher.FluxMap$MapConditionalSubscriber.onNext(FluxMap.java:224) at reactor.netty.channel.FluxReceive.drainReceiver(FluxReceive.java:292) at reactor.netty.channel.FluxReceive.onInboundNext(FluxReceive.java:401) at reactor.netty.channel.ChannelOperations.onInboundNext(ChannelOperations.java:411) at reactor.netty.channel.ChannelOperationsHandler.channelRead(ChannelOperationsHandler.java:113) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:336) at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:308) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:444) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1410) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:440) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:919) at io.netty.channel.epoll.AbstractEpollStreamChannel$EpollStreamUnsafe.epollInReady(AbstractEpollStreamChannel.java:800) at io.netty.channel.epoll.EpollEventLoop.processReady(EpollEventLoop.java:499) at io.netty.channel.epoll.EpollEventLoop.run(EpollEventLoop.java:397) at io.netty.util.concurrent.SingleThreadEventExecutor$4.run(SingleThreadEventExecutor.java:997) at io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74) at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) at java.base/java.lang.Thread.run(Thread.java:833)
复现说明
存在可复现该问题的测试代码,可验证问题场景。
已尝试的解决方案
- 切换
CoroutinesCrudRepository与ReactiveStreamsCrudRepository的不同组合 - 调整
@Transactional注解的使用方式,包括声明式事务配置
目前未找到明确错误规律,寻求可行解决方法。
内容的提问来源于stack exchange,提问作者Filipe Nascimento
相关产品推荐
相关产品推荐

