flink拉取kafka数据到doris报错版本不对齐报错

原因 1. 我的flink是1.7.2,hadoop是3.3.4,doris是3.1.0,streampark2.1.5,用的flink-doris-connector-1.17,版本原来是1.5.2;原来写入doris是2.1.11没问题,后面doris升级成3.1.0存算一体的作业起不来,但是线上用的selectdb 4.0.4没问题。

现象

查看日志是以下
xxxx是以下作业的Cluster Id

yarn logs -applicationId   xxxx

在这里插入图片描述

Send request to Doris FE ‘http://x.x.x.x:18030/api/bb/aa/_schema’ with user ‘root’.
2025-09-17 15:06:09,157 ERROR org.apache.doris.flink.table.DorisDynamicTableSink [] - Doris FE’s response cannot map to schema. res: {“keysType”:“UNIQUE_KEYS”,“properties”:


Send request to Doris FE 'http://x.x.x.x:18030/api/bb/aa/_schema' with user 'root'.
2025-09-17 15:06:09,157 ERROR org.apache.doris.flink.table.DorisDynamicTableSink           [] - Doris FE's response cannot map to schema. res: {"keysType":"UNIQUE_KEYS","properties":[{"name":"imei","aggregation_type":"","comment":"","is_nullable":"Yes","type":"BIGINT"},{"name":"time","aggregation_type":"","comment":"","is_nullable":"Yes","type":"BIGINT"},{"name":"dt","aggregation_type":"","comment":"","is_nullable":"Yes","type":"DATEV2"}],"status":200}
org.apache.doris.shaded.com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException: Unrecognized field "is_nullable" (class org.apache.doris.flink.rest.models.Field), not marked as ignorable (6 known properties: "aggregation_type", "type", "name", "precision", "comment", "scale"])
 at [Source: (String)"{"keysType":"UNIQUE_KEYS","properties":[{"name":"imei","aggregation_type":"","comment":"","is_nullable":"Yes","type":"BIGINT"},{"name":"time","aggregation_type":"","comment":"","is_nullable":"Yes","type":"BIGINT"},{"name":"dt","aggregation_type":"","comment":"","is_nullable":"Yes","type":"DATEV2"},{"name":"NTCs1",""[truncated 12808 chars]; line: 1, column: 106] (through reference chain: org.apache.doris.flink.rest.models.Schema["properties"]->java.util.ArrayList[0]->org.apache.doris.flink.rest.models.Field["is_nullable"])
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.exc.UnrecognizedPropertyException.from(UnrecognizedPropertyException.java:61) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.DeserializationContext.handleUnknownProperty(DeserializationContext.java:1127) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.std.StdDeserializer.handleUnknownProperty(StdDeserializer.java:2023) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.handleUnknownProperty(BeanDeserializerBase.java:1700) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializerBase.handleUnknownVanilla(BeanDeserializerBase.java:1678) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializer.vanillaDeserialize(BeanDeserializer.java:319) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserialize(BeanDeserializer.java:176) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.std.CollectionDeserializer._deserializeFromArray(CollectionDeserializer.java:355) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.std.CollectionDeserializer.deserialize(CollectionDeserializer.java:244) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.std.CollectionDeserializer.deserialize(CollectionDeserializer.java:28) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.impl.MethodProperty.deserializeAndSet(MethodProperty.java:129) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializer.vanillaDeserialize(BeanDeserializer.java:313) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.BeanDeserializer.deserialize(BeanDeserializer.java:176) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.deser.DefaultDeserializationContext.readRootValue(DefaultDeserializationContext.java:323) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.ObjectMapper._readMapAndClose(ObjectMapper.java:4674) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3629) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.shaded.com.fasterxml.jackson.databind.ObjectMapper.readValue(ObjectMapper.java:3597) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.flink.rest.RestService.parseSchema(RestService.java:545) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.flink.rest.RestService.getSchema(RestService.java:513) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.flink.rest.RestService.isUniqueKeyType(RestService.java:525) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.doris.flink.table.DorisDynamicTableSink.getSinkRuntimeProvider(DorisDynamicTableSink.java:84) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at org.apache.flink.table.planner.plan.nodes.exec.common.CommonExecSink.createSinkTransformation(CommonExecSink.java:158) ~[?:?]
        at org.apache.flink.table.planner.plan.nodes.exec.stream.StreamExecSink.translateToPlanInternal(StreamExecSink.java:176) ~[?:?]
        at org.apache.flink.table.planner.plan.nodes.exec.ExecNodeBase.translateToPlan(ExecNodeBase.java:161) ~[?:?]
        at org.apache.flink.table.planner.delegation.StreamPlanner.$anonfun$translateToPlan$1(StreamPlanner.scala:85) ~[?:?]
        at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:233) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.Iterator.foreach(Iterator.scala:937) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.Iterator.foreach$(Iterator.scala:937) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.AbstractIterator.foreach(Iterator.scala:1425) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.IterableLike.foreach(IterableLike.scala:70) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.IterableLike.foreach$(IterableLike.scala:69) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.AbstractIterable.foreach(Iterable.scala:54) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.TraversableLike.map(TraversableLike.scala:233) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.TraversableLike.map$(TraversableLike.scala:226) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at scala.collection.AbstractTraversable.map(Traversable.scala:104) ~[flink-scala_2.12-1.17.2.jar:1.17.2]
        at org.apache.flink.table.planner.delegation.StreamPlanner.translateToPlan(StreamPlanner.scala:84) ~[?:?]
        at org.apache.flink.table.planner.delegation.PlannerBase.translate(PlannerBase.scala:197) ~[?:?]
        at org.apache.flink.table.api.internal.TableEnvironmentImpl.translate(TableEnvironmentImpl.java:1805) ~[flink-table-api-java-uber-1.17.2.jar:1.17.2]
        at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:881) ~[flink-table-api-java-uber-1.17.2.jar:1.17.2]
        at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:991) ~[flink-table-api-java-uber-1.17.2.jar:1.17.2]
        at org.apache.flink.table.api.internal.TablePipelineImpl.execute(TablePipelineImpl.java:57) ~[flink-table-api-java-uber-1.17.2.jar:1.17.2]
        at org.apache.flink.table.api.Table.executeInsert(Table.java:1064) ~[flink-table-api-java-uber-1.17.2.jar:1.17.2]
        at com.yfzx.warehouse.realtime.GPS.odsdevice.Task.handle(Task.java:324) ~[kafka-to-doris_ods_device_basic_info-1.0-SNAPSHOT.jar:?]
        at com.yfzx.warehouse.realtime.common.base.BaseSqlApp.start(BaseSqlApp.java:41) ~[realtime-common-1.0-SNAPSHOT.jar:?]
        at com.yfzx.warehouse.realtime.GPS.odsdevice.Task.main(Task.java:10) ~[kafka-to-doris_ods_device_basic_info-1.0-SNAPSHOT.jar:?]
        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[?:1.8.0_212]
        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) ~[?:1.8.0_212]
        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[?:1.8.0_212]
        at java.lang.reflect.Method.invoke(Method.java:498) ~[?:1.8.0_212]
        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:355) ~[flink-dist-1.17.2.jar:1.17.2]
        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222) ~[flink-dist-1.17.2.jar:1.17.2]
        at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:105) ~[flink-dist-1.17.2.jar:1.17.2]
        at org.apache.flink.client.deployment.application.ApplicationDispatcherBootstrap.runApplicationEntryPoint(ApplicationDispatcherBootstrap.java:301) ~[flink-dist-1.17.2.jar:1.17.2]
        at org.apache.flink.client.deployment.application.ApplicationDispatcherBootstrap.lambda$runApplicationAsync$2(ApplicationDispatcherBootstrap.java:254) ~[flink-dist-1.17.2.jar:1.17.2]
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[?:1.8.0_212]
        at java.util.concurrent.FutureTask.run(FutureTask.java:266) ~[?:1.8.0_212]
        at org.apache.flink.runtime.concurrent.akka.ActorSystemScheduledExecutorAdapter$ScheduledFutureTask.run(ActorSystemScheduledExecutorAdapter.java:171) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.lambda$withContextClassLoader$0(ClassLoadingUtils.java:41) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:49) [flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(ForkJoinExecutorConfigurator.scala:48) [flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:157) [?:1.8.0_212]
2025-09-17 15:06:09,181 WARN  org.apache.flink.client.deployment.application.ApplicationDispatcherBootstrap [] - Application failed unexpectedly: 
java.util.concurrent.CompletionException: org.apache.flink.client.deployment.application.ApplicationExecutionException: Could not execute application.
        at java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:292) ~[?:1.8.0_212]
        at java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:308) ~[?:1.8.0_212]
        at java.util.concurrent.CompletableFuture.uniCompose(CompletableFuture.java:943) ~[?:1.8.0_212]
        at java.util.concurrent.CompletableFuture$UniCompose.tryFire(CompletableFuture.java:926) ~[?:1.8.0_212]
        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474) ~[?:1.8.0_212]
        at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1977) ~[?:1.8.0_212]
        at org.apache.flink.client.deployment.application.ApplicationDispatcherBootstrap.runApplicationEntryPoint(ApplicationDispatcherBootstrap.java:337) ~[flink-dist-1.17.2.jar:1.17.2]
        at org.apache.flink.client.deployment.application.ApplicationDispatcherBootstrap.lambda$runApplicationAsync$2(ApplicationDispatcherBootstrap.java:254) ~[flink-dist-1.17.2.jar:1.17.2]
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[?:1.8.0_212]
        at java.util.concurrent.FutureTask.run(FutureTask.java:266) ~[?:1.8.0_212]
        at org.apache.flink.runtime.concurrent.akka.ActorSystemScheduledExecutorAdapter$ScheduledFutureTask.run(ActorSystemScheduledExecutorAdapter.java:171) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.lambda$withContextClassLoader$0(ClassLoadingUtils.java:41) ~[flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:49) [flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(ForkJoinExecutorConfigurator.scala:48) [flink-rpc-akka_35d12d54-01dd-4e6e-9833-4540f0a3a6bb.jar:1.17.2]
        at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692) [?:1.8.0_212]
        at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:157) [?:1.8.0_212]
Caused by: org.apache.flink.client.deployment.application.ApplicationExecutionException: Could not execute application.
        ... 13 more

解决办法

把maven的pom升级flink-doris-connector


 <dependency>
                <groupId>org.apache.doris</groupId>
                <artifactId>flink-doris-connector-1.17</artifactId>
                <version>1.5.2</version>
</dependency>

 <dependency>
                <groupId>org.apache.doris</groupId>
                <artifactId>flink-doris-connector-1.17</artifactId>
                <version>25.1.0</version>
</dependency>

重启打包解决问题

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐