Kotlin+Spring中Reactive Cassandra Driver的int类型参数异常解决
解决Kotlin访问Cassandra时的
Primitive type 'int' used as type parameter错误 问题场景
在Kotlin编写的Spring应用中,使用ReactiveCqlTemplate访问Cassandra时触发Primitive type 'int' used as type parameter异常,错误出现在row.getList("c1", Int::class.java)代码处。由于业务场景中Schema会频繁动态变动,无法使用CassandraTemplate。
原代码片段:
@Component class Tst( private val cql: ReactiveCqlTemplate, ) { val ks = "test1" val table = "some_table" init { runBlocking { cql.execute("create keyspace if not exists $ks with replication = {'class': 'SimpleStrategy', 'replication_factor': 1};") .awaitSingle() cql.execute("drop table if exists $ks.$table;").awaitSingle() cql.execute( """create table if not exists $ks.$table ( id bigint, c1 list<int>, primary key (id) );""" ).awaitSingle() val someList = listOf(1, 2, 4, 67, 8, 5, 3, 3, 79, 7453457) cql.execute( "insert into $ks.$table (id, c1) values (?,?)", 1L, someList ).awaitSingle() val result = cql.query("select c1 from $ks.$table where id = ${1L}") { row, _ -> val result = row.getList("c1", Int::class.java) // 错误行 result }.next().awaitSingle() println(result) } } }
完整异常栈:
Caused by: java.lang.IllegalArgumentException: Primitive type 'int' used as type parameter at com.datastax.oss.driver.shaded.guava.common.base.Preconditions.checkArgument(Preconditions.java:440) at com.datastax.oss.driver.shaded.guava.common.reflect.Types.disallowPrimitiveType(Types.java:525) at com.datastax.oss.driver.shaded.guava.common.reflect.Types.access$200(Types.java:54) at com.datastax.oss.driver.shaded.guava.common.reflect.Types$ParameterizedTypeImpl.<init>(Types.java:267) at com.datastax.oss.driver.shaded.guava.common.reflect.Types.newParameterizedType(Types.java:102) at com.datastax.oss.driver.shaded.guava.common.reflect.Types.newParameterizedTypeWithOwner(Types.java:91) at com.datastax.oss.driver.shaded.guava.common.reflect.TypeResolver.resolveParameterizedType(TypeResolver.java:265) at com.datastax.oss.driver.shaded.guava.common.reflect.TypeResolver.resolveType(TypeResolver.java:220) at com.datastax.oss.driver.shaded.guava.common.reflect.TypeToken.where(TypeToken.java:231) at com.datastax.oss.driver.api.core.type.reflect.GenericType.listOf(GenericType.java:126) at com.datastax.oss.driver.api.core.data.GettableByIndex.getList(GettableByIndex.java:499) at com.datastax.oss.driver.api.core.data.GettableByName.getList(GettableByName.java:580) at com.example.Tst$1.invokeSuspend$lambda$0(Tst.kt:49) at org.springframework.data.cassandra.core.cql.ReactiveRowMapperResultSetExtractor.lambda$extractData$0(ReactiveRowMapperResultSetExtractor.java:61) at reactor.core.publisher.FluxHandleFuseable$HandleFuseableSubscriber.onNext(FluxHandleFuseable.java:179) at reactor.core.publisher.FluxFlattenIterable$FlattenIterableSubscriber.drainAsync(FluxFlattenIterable.java:453) at reactor.core.publisher.FluxFlattenIterable$FlattenIterableSubscriber.drain(FluxFlattenIterable.java:724) at reactor.core.publisher.FluxFlattenIterable$FlattenIterableSubscriber.onNext(FluxFlattenIterable.java:256) at reactor.core.publisher.FluxExpand$ExpandBreathSubscriber.onNext(FluxExpand.java:118) at reactor.core.publisher.Operators$ScalarSubscription.request(Operators.java:2571) at reactor.core.publisher.Operators$MultiSubscriptionSubscriber.set(Operators.java:2367) at reactor.core.publisher.FluxExpand$ExpandBreathSubscriber.onSubscribe(FluxExpand.java:112) at reactor.core.publisher.MonoJust.subscribe(MonoJust.java:55) at reactor.core.publisher.Mono.subscribe(Mono.java:4563) at reactor.core.publisher.FluxExpand$ExpandBreathSubscriber.drainQueue(FluxExpand.java:179) at reactor.core.publisher.MonoExpand.subscribeOrReturn(MonoExpand.java:54) at reactor.core.publisher.Flux.subscribe(Flux.java:8821) at reactor.core.publisher.MonoFlatMapMany$FlatMapManyMain.onNext(MonoFlatMapMany.java:196) at reactor.core.publisher.FluxMap$MapSubscriber.onNext(FluxMap.java:122) at reactor.core.publisher.MonoCompletionStage$MonoCompletionStageSubscription.apply(MonoCompletionStage.java:121) at reactor.core.publisher.MonoCompletionStage$MonoCompletionStageSubscription.apply(MonoCompletionStage.java:67) at java.base/java.util.concurrent.CompletableFuture.uniHandle(CompletableFuture.java:934) at java.base/java.util.concurrent.CompletableFuture$UniHandle.tryFire(CompletableFuture.java:911) at java.base/java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:510) at java.base/java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:2147) at com.datastax.oss.driver.internal.core.cql.CqlRequestHandler.setFinalResult(CqlRequestHandler.java:324) at com.datastax.oss.driver.internal.core.cql.CqlRequestHandler.access$1500(CqlRequestHandler.java:95) at com.datastax.oss.driver.internal.core.cql.CqlRequestHandler$NodeResponseCallback.onResponse(CqlRequestHandler.java:655) at com.datastax.oss.driver.internal.core.channel.InFlightHandler.channelRead(InFlightHandler.java:257) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:442) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) at io.netty.handler.timeout.IdleStateHandler.channelRead(IdleStateHandler.java:289) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:442) at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:420) at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:412) at io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:103) 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:346) at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:318) 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.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:166) at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:788) at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:724) at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:650) at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:562) 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:840)
错误原因
Cassandra Java驱动依赖的Guava库在处理泛型类型时,明确禁止使用JVM原生类型(如int)作为泛型参数。而Kotlin中Int::class.java返回的是JVM原生int的Class对象,不是装箱后的Integer类型,因此触发该校验错误。
解决方案
将获取类型的方式替换为装箱类型的Class对象,有两种可行方式:
方式一:直接使用Integer::class.java
修改错误行代码为:
val result = row.getList("c1", Integer::class.java)
方式二:使用Kotlin的javaObjectType获取装箱类型
利用Kotlin反射API中的javaObjectType属性,直接获取对应装箱类型的Class对象:
val result = row.getList("c1", Int::class.javaObjectType)
修改后的完整代码片段
@Component class Tst( private val cql: ReactiveCqlTemplate, ) { val ks = "test1" val table = "some_table" init { runBlocking { cql.execute("create keyspace if not exists $ks with replication = {'class': 'SimpleStrategy', 'replication_factor': 1};") .awaitSingle() cql.execute("drop table if exists $ks.$table;").awaitSingle() cql.execute( """create table if not exists $ks.$table ( id bigint, c1 list<int>, primary key (id) );""" ).awaitSingle() val someList = listOf(1, 2, 4, 67, 8, 5, 3, 3, 79, 7453457) cql.execute( "insert into $ks.$table (id, c1) values (?,?)", 1L, someList ).awaitSingle() val result = cql.query("select c1 from $ks.$table where id = ${1L}") { row, _ -> val result = row.getList("c1", Int::class.javaObjectType) // 修改后的代码 result }.next().awaitSingle() println(result) } } }
内容的提问来源于stack exchange,提问作者agathis
相关产品推荐
相关产品推荐

