• [问题求助] 【大数据产品】【Flink功能】ambari部署大数据平台缺失Flink组件
    【功能模块】【操作步骤&问题现象】ambari部署大数据平台缺失Flink组件,源是hdp3.1,有没编译好的Flink安装包【截图信息】【日志信息】(可选,上传日志内容或者附件)
  • [问题求助] flink无法提交job作业
    【功能模块】【操作步骤&问题现象】1、将flinkkafkaexample样例程序打包运行,报无法提交作业的错误。2、【截图信息】【日志信息】(可选,上传日志内容或者附件) The program finished with the following exception:org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: java.util.concurrent.ExecutionException: org.apache.flink.runtime.client.JobSubmissionException: Failed to submit JobGraph.        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:335)        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:205)        at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:138)        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:700)        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:219)        at org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:932)        at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1005)        at java.security.AccessController.doPrivileged(Native Method)        at javax.security.auth.Subject.doAs(Subject.java:422)        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1737)        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1005)Caused by: java.lang.RuntimeException: java.util.concurrent.ExecutionException: org.apache.flink.runtime.client.JobSubmissionException: Failed to submit JobGraph.        at org.apache.flink.util.ExceptionUtils.rethrow(ExceptionUtils.java:199)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.executeAsync(StreamExecutionEnvironment.java:1772)        at org.apache.flink.streaming.api.environment.StreamContextEnvironment.executeAsync(StreamContextEnvironment.java:94)        at org.apache.flink.streaming.api.environment.StreamContextEnvironment.execute(StreamContextEnvironment.java:63)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1651)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1633)        at com.huawei.bigdata.flink.examples.WriteIntoKafka.main(WriteIntoKafka.java:48)        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)        at java.lang.reflect.Method.invoke(Method.java:498)        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:321)        ... 11 moreCaused by: java.util.concurrent.ExecutionException: org.apache.flink.runtime.client.JobSubmissionException: Failed to submit JobGraph.        at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)        at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.executeAsync(StreamExecutionEnvironment.java:1767)        ... 21 moreCaused by: org.apache.flink.runtime.client.JobSubmissionException: Failed to submit JobGraph.        at org.apache.flink.client.program.rest.RestClusterClient.lambda$submitJob$7(RestClusterClient.java:359)        at java.util.concurrent.CompletableFuture.uniExceptionally(CompletableFuture.java:884)        at java.util.concurrent.CompletableFuture$UniExceptionally.tryFire(CompletableFuture.java:866)        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)        at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1990)        at org.apache.flink.runtime.concurrent.FutureUtils.lambda$retryOperationWithDelay$8(FutureUtils.java:287)        at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774)        at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750)        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)        at java.util.concurrent.CompletableFuture.completeExceptionally(CompletableFuture.java:1990)        at org.apache.flink.runtime.rest.RestClient$ClientHandler.channelInactive(RestClient.java:483)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236)        at org.apache.flink.shaded.netty4.io.netty.channel.ChannelInboundHandlerAdapter.channelInactive(ChannelInboundHandlerAdapter.java:81)        at org.apache.flink.shaded.netty4.io.netty.handler.timeout.IdleStateHandler.channelInactive(IdleStateHandler.java:278)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236)        at org.apache.flink.shaded.netty4.io.netty.handler.stream.ChunkedWriteHandler.channelInactive(ChunkedWriteHandler.java:141)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236)        at org.apache.flink.shaded.netty4.io.netty.channel.ChannelInboundHandlerAdapter.channelInactive(ChannelInboundHandlerAdapter.java:81)        at org.apache.flink.shaded.netty4.io.netty.handler.codec.MessageAggregator.channelInactive(MessageAggregator.java:438)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236)        at org.apache.flink.shaded.netty4.io.netty.channel.CombinedChannelDuplexHandler$DelegatingChannelHandlerContext.fireChannelInactive(CombinedChannelDuplexHandler.java:420)        at org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.channelInputClosed(ByteToMessageDecoder.java:393)        at org.apache.flink.shaded.netty4.io.netty.handler.codec.ByteToMessageDecoder.channelInactive(ByteToMessageDecoder.java:358)        at org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpClientCodec$Decoder.channelInactive(HttpClientCodec.java:282)        at org.apache.flink.shaded.netty4.io.netty.channel.CombinedChannelDuplexHandler.channelInactive(CombinedChannelDuplexHandler.java:223)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.fireChannelInactive(AbstractChannelHandlerContext.java:236)        at org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline$HeadContext.channelInactive(DefaultChannelPipeline.java:1416)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:257)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannelHandlerContext.invokeChannelInactive(AbstractChannelHandlerContext.java:243)        at org.apache.flink.shaded.netty4.io.netty.channel.DefaultChannelPipeline.fireChannelInactive(DefaultChannelPipeline.java:912)        at org.apache.flink.shaded.netty4.io.netty.channel.AbstractChannel$AbstractUnsafe$8.run(AbstractChannel.java:816)        at org.apache.flink.shaded.netty4.io.netty.util.concurrent.AbstractEventExecutor.safeExecute(AbstractEventExecutor.java:163)        at org.apache.flink.shaded.netty4.io.netty.util.concurrent.SingleThreadEventExecutor.runAllTasks(SingleThreadEventExecutor.java:416)        at org.apache.flink.shaded.netty4.io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:515)        at org.apache.flink.shaded.netty4.io.netty.util.concurrent.SingleThreadEventExecutor$5.run(SingleThreadEventExecutor.java:918)        at org.apache.flink.shaded.netty4.io.netty.util.internal.ThreadExecutorMap$2.run(ThreadExecutorMap.java:74)        at java.lang.Thread.run(Thread.java:748)Caused by: org.apache.flink.runtime.concurrent.FutureUtils$RetryException: Could not complete the operation. Number of retries has been exhausted.        at org.apache.flink.runtime.concurrent.FutureUtils.lambda$retryOperationWithDelay$8(FutureUtils.java:284)        ... 41 moreCaused by: java.util.concurrent.CompletionException: org.apache.flink.runtime.rest.ConnectionClosedException: Channel became inactive.        at java.util.concurrent.CompletableFuture.encodeRelay(CompletableFuture.java:326)        at java.util.concurrent.CompletableFuture.completeRelay(CompletableFuture.java:338)        at java.util.concurrent.CompletableFuture.uniRelay(CompletableFuture.java:925)        at java.util.concurrent.CompletableFuture$UniRelay.tryFire(CompletableFuture.java:913)        ... 39 moreCaused by: org.apache.flink.runtime.rest.ConnectionClosedException: Channel became inactive.        ... 37 more
  • [其他问题] 在flink开发指南中,flinkkafkaexaple打包在linux上运行报错,报的错是:提取元数据时,kafka超时已过期
    【功能模块】【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)报错信息如下: The program finished with the following exception:org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: 9e3c78adac04823062e5328231935728)        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:335)        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:205)        at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:138)        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:700)        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:219)        at org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:932)        at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1005)        at java.security.AccessController.doPrivileged(Native Method)        at javax.security.auth.Subject.doAs(Subject.java:422)        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1737)        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1005)Caused by: java.util.concurrent.ExecutionException: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: 9e3c78adac04823062e5328231935728)        at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)        at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908)        at org.apache.flink.streaming.api.environment.StreamContextEnvironment.execute(StreamContextEnvironment.java:83)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1651)        at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1633)        at com.huawei.bigdata.flink.examples.ReadFromKafka.main(ReadFromKafka.java:59)        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)        at java.lang.reflect.Method.invoke(Method.java:498)        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:321)        ... 11 moreCaused by: org.apache.flink.client.program.ProgramInvocationException: Job failed (JobID: 9e3c78adac04823062e5328231935728)        at org.apache.flink.client.deployment.ClusterClientJobClientAdapter.lambda$null$6(ClusterClientJobClientAdapter.java:112)        at java.util.concurrent.CompletableFuture.uniApply(CompletableFuture.java:616)        at java.util.concurrent.CompletableFuture$UniApply.tryFire(CompletableFuture.java:591)        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)        at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1975)        at org.apache.flink.client.program.rest.RestClusterClient.lambda$pollResourceAsync$21(RestClusterClient.java:565)        at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774)        at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750)        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)        at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1975)        at org.apache.flink.runtime.concurrent.FutureUtils.lambda$retryOperationWithDelay$8(FutureUtils.java:291)        at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:774)        at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:750)        at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:488)        at java.util.concurrent.CompletableFuture.postFire(CompletableFuture.java:575)        at java.util.concurrent.CompletableFuture$UniCompose.tryFire(CompletableFuture.java:943)        at java.util.concurrent.CompletableFuture$Completion.run(CompletableFuture.java:456)        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)        at java.lang.Thread.run(Thread.java:748)Caused by: org.apache.flink.runtime.client.JobExecutionException: Job execution failed.        at org.apache.flink.runtime.jobmaster.JobResult.toJobExecutionResult(JobResult.java:147)        at org.apache.flink.client.deployment.ClusterClientJobClientAdapter.lambda$null$6(ClusterClientJobClientAdapter.java:110)        ... 19 moreCaused by: org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy        at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:110)        at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:76)        at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:192)        at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:186)        at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:180)        at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:484)        at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:380)        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)        at java.lang.reflect.Method.invoke(Method.java:498)        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:279)        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:194)        at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:152)        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26)        at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21)        at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:123)        at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:170)        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)        at akka.actor.Actor$class.aroundReceive(Actor.scala:517)        at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225)        at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592)        at akka.actor.ActorCell.invoke(ActorCell.scala:561)        at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258)        at akka.dispatch.Mailbox.run(Mailbox.scala:225)        at akka.dispatch.Mailbox.exec(Mailbox.scala:235)        at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)        at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)        at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)        at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)Caused by: org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata
  • [基础组件] 华为FI C80 kerberos认证环境下,Flink 无法提交作业到yarn上,既是flink配置了认证文件也不行,求助
    因公司需要,需要在华为FIC80环境下执行Flink作业,但是目前遇到了问题,过程如下:1、先将Flink 1.12.0 版本解压后配置conf/flink-conf.yml 的安全认证,如下:security.kerberos.login.use-ticket-cache: truesecurity.kerberos.login.keytab: /home/kgp/user.keytabsecurity.kerberos.login.principal: kgp@HADOOP.COMsecurity.kerberos.login.contexts: Client,KafkaClient环境变量中的hadoop_home 和 hadoop_classpath都配置了,flink是可以找到yarn的2、然后尝试执行官网的demo案例:先执行了kinit命令:kinit -kt user.keytab kgp然后提交作业./bin/flink run -t yarn-per-job --detached ./examples/streaming/TopSpeedWindowing.jar3、结果出现异常:此时yarn UI 界面没有此任务,按照提示执行yarn命令来查看日志命令:yarn logs -applicationId application_1616467846343_00314、执行后得到的日志为:所用到的票据皆是从FI管理界面上下载的,可以保证正确求助
  • [性能调优] 【flink checkpoint】checkpoint超时导致失败
    【功能模块】消费kafka消息写hdfs,flink 1.10,FI8.0.2【操作步骤&问题现象】1、操作步骤source-sink在一个task,并行度240,FsStateBackend1.1 RollingPolicy:batchSize=1024000000rolloverInterval=6000001.2 CheckpointConfig:chkInterval=600000chkMinPause=30000chkTimeout=500000chkMaxConcurrentNum=1FsStateBackend1.3 数据量大小 1.7T/10mins1.4 分配的yarn队列资源,内存、CPU均未用完2、问题现象2.1 怎么排查checkpoint为什么慢2.2 checkpoint sync和async分别是做什么的【截图信息】【日志信息】(可选,上传日志内容或者附件)对应时间点没做checkpoint的subtask的tm的log日志已保存,因为一些关键信息就不贴在这了。希望有大神指导!!!
  • [问题求助] 【fi 华为flink sql】【XXX功能】flink sql 是怎么读和写Gauss高斯数据库的?
    【功能模块】【fi  华为flink  sql】【XXX功能】flink sql 是怎么读和写Gauss高斯数据库的?【操作步骤&问题现象】1、我想迁移这部分功能到开源的flink 1.12上, 使得  开源的flink  sql也支持  读和写 Gauss高斯数据库 ?2、【截图信息】【日志信息】(可选,上传日志内容或者附件)
  • [问题求助] 【flink sql】【 尝试 使用 huaweicloud-mrs-example-mrs-2.1】调flinksql???
    【功能模块】【flink sql】【 尝试 使用 huaweicloud-mrs-example-mrs-2.1】调flinksql???【操作步骤&问题现象】1、基于 huaweicloud-mrs-example-mrs-2.1 在 flink 1.7.2版本的 环境 运行 flink  sql报错了?? 2、 测试 代码 package com.huawei.flink.example.sqljoin; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.TableEnvironment; import org.apache.flink.table.api.java.StreamTableEnvironment; //flink run -c   com.huawei.flink.example.sqljoin.SqlSim   -m yarn-cluster    /srv/BigData/data1/task_script_dir/onedata_ec_analyse/flink/s5.jar public class SqlSim {     public static final String KAFKA_TABLE_SOURCE_DDL = "" +             "CREATE SOURCE STREAM stream_id (attr_name attr_type (',' attr_name attr_type)* )\n" +             "  WITH (\n" +             "    type = \"kafka\",\n" +             "    kafka_bootstrap_servers = \"\",\n" +             "    kafka_group_id = \"\",\n" +             "    kafka_topic = \"\",\n" +             "    encode = \"json\"\n" +             "  )\n" +             "  (TIMESTAMP BY timeindicator (',' timeindicator)?);timeindicator:PROCTIME '.' PROCTIME| ID '.' ROWTIME";     public static void main(String[] args) throws Exception{         System.err.println("--77777*--SqlSim--- ");         StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();         StreamTableEnvironment tEnv = TableEnvironment.getTableEnvironment(env); //        ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); //        BatchTableEnvironment tableEnv = BatchTableEnvironment.create(env);         System.err.println("--35 *--SqlSim--- ");         tEnv.sqlUpdate(KAFKA_TABLE_SOURCE_DDL);         System.err.println("--38  38*--SqlSim--- ");         env.execute();         System.err.println("--42 42*--SqlSim--- ");     } }测试官网的例子" CREATE SOURCE STREAM stream_id (attr_name attr_type (',' attr_name attr_type)* )   WITH (     type = "kafka",     kafka_bootstrap_servers = "",     kafka_group_id = "",     kafka_topic = "",     encode = "json"   )   (TIMESTAMP BY timeindicator (',' timeindicator)?);timeindicator:PROCTIME '.' PROCTIME| ID '.' ROWTIME【截图信息】【日志信息】(可选,上传日志内容或者附件)apache.flink.client.program.ProgramInvocationException: The main method caused an error.        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:546)        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)        at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:430)        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:814)        at org.apache.flink.client.cli.CliFrontend.runProgram(CliFrontend.java:288)        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:213)        at org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:1051)        at org.apache.flink.client.cli.CliFrontend.lambda$main$11(CliFrontend.java:1127)        at java.security.AccessController.doPrivileged(Native Method)        at javax.security.auth.Subject.doAs(Subject.java:422)        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1729)        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1127)Caused by: org.apache.flink.table.api.SqlParserException: SQL parse failed. Encountered "CREATE" at line 1, column 1.Was expecting one of:    "SET" ...    "RESET" ...    "ALTER" ...    "WITH" ...    "+" ...    "-" ...    "NOT" ...    "EXISTS" ...    <UNSIGNED_INTEGER_LITERAL> ...    <DECIMAL_NUMERIC_LITERAL> ...    <APPROX_NUMERIC_LITERAL> ...    <BINARY_STRING_LITERAL> ...    <PREFIXED_STRING_LITERAL> ...    <QUOTED_STRING> ...    <UNICODE_STRING_LITERAL> ...    "TRUE" ...    "FALSE" ...    "UNKNOWN" ...    "NULL" ...    <LBRACE_D> ...    <LBRACE_T> ...    <LBRACE_TS> ...    "DATE" ...    "TIME" ...
  • [问题求助] 【flink 1.7.2】【向Kafka生产】ClassNotFoundException: org.apache.kafka.
    【功能模块】bin/flink run --class com.huawei.flink.example.sqljoin.WriteIntoKafka4SQLJoin /opt/Flink_test/flink-examples-1.0.jar --topic topic-test --bootstrap.servers xxx.xxx.xxx.xxx:21005运行测试 例子   com.huawei.flink.example.sqljoin.WriteIntoKafka4SQLJoin运行不通过  【操作步骤&问题现象】1、2、【截图信息】【日志信息】(可选,上传日志内容或者附件)--sasl.kerberos.service.name kafka --ssl.truststore.location /home/truststore.jks --ssl.truststore.password huawei******************************************************************************************<topic> is the kafka topic name<bootstrap.servers> is the ip:port list of brokers******************************************************************************************java.lang.NoClassDefFoundError: org/apache/kafka/clients/producer/Callback        at com.huawei.flink.example.sqljoin.WriteIntoKafka4SQLJoin.main(WriteIntoKafka4SQLJoin.java:53)        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)        at java.lang.reflect.Method.invoke(Method.java:498)        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:529)        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)        at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:430)        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:814)        at org.apache.flink.client.cli.CliFrontend.runProgram(CliFrontend.java:288)        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:213)        at org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:1051)        at org.apache.flink.client.cli.CliFrontend.lambda$main$11(CliFrontend.java:1127)        at java.security.AccessController.doPrivileged(Native Method)        at javax.security.auth.Subject.doAs(Subject.java:422)        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1729)        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1127)Caused by: java.lang.ClassNotFoundException: org.apache.kafka.clients.producer.Callback        at java.net.URLClassLoader.findClass(URLClassLoader.java:382)        at java.lang.ClassLoader.loadClass(ClassLoader.java:424)        at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349)        at java.lang.ClassLoader.loadClass(ClassLoader.java:357)        ... 18 more
  • [问题求助] 想问一下 华为版 fi 的flink 重点构建的 代码 DataStream,有没有开源GitHub , 有分支开源可以看的
    想问一下  华为版 fi 的flink  重点构建的 代码 ,有没有开源GitHub , 有分支开源可以看的代码吗? Flink在当前版本中重点构建如下特性:DataStreamCheckpoint窗口Job Pipeline配置表其他特性继承开源社区,不做增强,具体请参 :https://ci.apache.org/projects/flink/flink-docs-release-1.12/。
  • [问题求助] 我们要调 flink sql 例子, huaweicloud-mrs-example有好几个分支,怎么知道用哪个分支?
    【功能模块】我们要调 flink sql 例子, huaweicloud-mrs-example有好几个分支,怎么知道用哪个分支?【操作步骤&问题现象】1、 导入了   https://github.com/huaweicloud/huaweicloud-mrs-example/branches   的   mrs-3.0.2Updated last mont   分支 Compar2、   测试     flink  run  -m  yarn-cluster   --class com.huawei.bigdata.flink.examples.JavaStreamSqlExample   ./../tmp/s.jar测试  com.huawei.bigdata.flink.examples.JavaStreamSqlExample的 例子,不通过  【截图信息】【日志信息】(可选,上传日志内容或者附件)bin/flink run --class com.huawei.bigdata.flink.examples.JavaStreamSqlExample <path of StreamSqlExample jar>*********************************************************************************java.lang.NoClassDefFoundError: org/apache/flink/table/factories/StreamTableSourceFactory        at java.lang.ClassLoader.defineClass1(Native Method)        at java.lang.ClassLoader.defineClass(ClassLoader.java:763)        at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:142)        at java.net.URLClassLoader.defineClass(URLClassLoader.java:468)        at java.net.URLClassLoader.access$100(URLClassLoader.java:74)        at java.net.URLClassLoader$1.run(URLClassLoader.java:369)        at java.net.URLClassLoader$1.run(URLClassLoader.java:363)        at java.security.AccessController.doPrivileged(Native Method)        at java.net.URLClassLoader.findClass(URLClassLoader.java:362)        at java.lang.ClassLoader.loadClass(ClassLoader.java:424)        at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349)        at java.lang.ClassLoader.loadClass(ClassLoader.java:357)        at java.lang.ClassLoader.defineClass1(Native Method)        at java.lang.ClassLoader.defineClass(ClassLoader.java:763)        at java.security.SecureClassLoader.defineClass(SecureClassLoader.java:142)        at java.net.URLClassLoader.defineClass(URLClassLoader.java:468)        at java.net.URLClassLoader.access$100(URLClassLoader.java:74)        at java.net.URLClassLoader$1.run(URLClassLoader.java:369)        at java.net.URLClassLoader$1.run(URLClassLoader.java:363)        at java.security.AccessController.doPrivileged(Native Method)        at java.net.URLClassLoader.findClass(URLClassLoader.java:362)        at java.lang.ClassLoader.loadClass(ClassLoader.java:424)        at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349)        at java.lang.ClassLoader.loadClass(ClassLoader.java:411)        at java.lang.ClassLoader.loadClass(ClassLoader.java:357)        at java.lang.Class.forName0(Native Method)        at java.lang.Class.forName(Class.java:348)        at java.util.ServiceLoader$LazyIterator.nextService(ServiceLoader.java:370)        at java.util.ServiceLoader$LazyIterator.next(ServiceLoader.java:404)        at java.util.ServiceLoader$1.next(ServiceLoader.java:480)        at java.util.Iterator.forEachRemaining(Iterator.java:116)        at org.apache.flink.table.factories.TableFactoryService.discoverFactories(TableFactoryService.java:214)        at org.apache.flink.table.factories.TableFactoryService.findAllInternal(TableFactoryService.java:170)        at org.apache.flink.table.factories.TableFactoryService.findAll(TableFactoryService.java:125)        at org.apache.flink.table.factories.ComponentFactoryService.find(ComponentFactoryService.java:48)        at org.apache.flink.table.api.java.internal.StreamTableEnvironmentImpl.lookupExecutor(StreamTableEnvironmentImpl.java:143)        at org.apache.flink.table.api.java.internal.StreamTableEnvironmentImpl.create(StreamTableEnvironmentImpl.java:118)        at org.apache.flink.table.api.java.StreamTableEnvironment.create(StreamTableEnvironment.java:112)        at com.huawei.bigdata.flink.examples.JavaStreamSqlExample.main(JavaStreamSqlExample.java:32)        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)        at java.lang.reflect.Method.invoke(Method.java:498)        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:529)        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)        at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:430)        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:814)        at org.apache.flink.client.cli.CliFrontend.runProgram(CliFrontend.java:288)        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:213)        at org.apache.flink.client.cli.CliFrontend.parseParameters(CliFrontend.java:1051)        at org.apache.flink.client.cli.CliFrontend.lambda$main$11(CliFrontend.java:1127)        at java.security.AccessController.doPrivileged(Native Method)        at javax.security.auth.Subject.doAs(Subject.java:422)        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1729)        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1127)Caused by: java.lang.ClassNotFoundException: org.apache.flink.table.factories.StreamTableSourceFactory        at java.net.URLClassLoader.findClass(URLClassLoader.java:382)        at java.lang.ClassLoader.loadClass(ClassLoader.java:424)        at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:349)        at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
  • [问题求助] FI 的 flink 1.7版本的 flink sql 推荐使用什么方法, 推荐用哪个工程做StreamSQLExampl
    FI  的 flink 1.7版本的 flink  sql  推荐使用什么方法, 推荐用哪个工程做StreamSQLExample.,是否直接从 GitHub 官网的 flink 下载 源码, 直接整 StreamSQLExample  ???
  • [问题求助] 编译flink 源码 打算,运行 StreamSQLExample ,不通过
    编译flink 源码 打算,运行 StreamSQLExample ,不通过JPS incremental annotation processing is disabled. Compilation results on partial recompilation may be inaccurate. Use build process "jps.track.ap.dependencies" VM flag to enable/disable incremental annotation processing environmentjava: 警告: 源发行版 11 需要目标发行版 11
  • [问题求助] Flink SQL&gt; SELECT &apos;Hello World&apos;; 最简单的 sql测试都不行?
    Flink SQL> SELECT 'Hello World';>[ERROR] Could not execute SQL statement. Reason:java.lang.NoSuchMethodError: scala.collection.JavaConversions$.deprecated$u0020seqAsJavaList(Lscala/collection/Seq;)Ljava/util/List;
  • [问题求助] DLI flink中可以写解密过程对数据解密么 怎么操作
    数据直接接入的dis,是加密好的,再用dli flink去处理dis的数据时怎么写入解密过程 求一个详细的解决方案。
  • [问题求助] 调fi flink 的 ,sql-client.sh 报错了 NoSuchFieldError: PYFILES_OPTIO
    添加了jar包之后,还是报错了  NoSuchFieldError: PYFILES_OPTION/FusionInsight/client/HDFS/hadoop/share/hadoop/common/lib/slf4j-log4j12-1.7.25.jar!/org/slf4j/impl/StaticLoggerBinder.class]SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation.SLF4J: Actual binding is of type [ch.qos.logback.classic.util.ContextSelectorStaticBinder]Exception in thread "main" java.lang.NoSuchFieldError: PYFILES_OPTION        at org.apache.flink.table.client.cli.CliOptionsParser.getEmbeddedModeClientOptions(CliOptionsParser.java:156)        at org.apache.flink.table.client.cli.CliOptionsParser.<clinit>(CliOptionsParser.java:139)        at org.apache.flink.table.client.SqlClient.main(SqlClient.java:195)
总条数:123 到第
上滑加载中