使用flinkkinesis producer连接到本地堆栈kinesis失败

aelbi1ox  于 2021-06-21  发布在  Flink
关注(0)|答案(1)|浏览(439)

我有https://github.com/localstack/localstack 在我的地方跑来做一些运动训练。该服务通过 localhost:4568 我可以使用awscli与之交互。
比如当我运行这个

  1. AWS_ACCESS_KEY_ID=x AWS_SECRET_ACCESS_KEY=x aws --region us-east-1 --endpoint-url http://localhost:4568/ kinesis create-stream --stream-name myStreamName --shard-count 1 --no-verify-ssl
  2. AWS_ACCESS_KEY_ID=x AWS_SECRET_ACCESS_KEY=x aws --region us-east-1 --endpoint-url http://localhost:4568/ kinesis list-streams

我明白了

  1. {
  2. "StreamNames": [
  3. "myStreamName"
  4. ]
  5. }

当我配置 FlinkKinesisProducer 但是像这样

  1. Properties kinesisConfig = new Properties();
  2. // Required configs
  3. kinesisConfig.put(AWSConfigConstants.AWS_REGION, "us-east-1");
  4. kinesisConfig.put(AWSConfigConstants.AWS_ACCESS_KEY_ID, "x");
  5. kinesisConfig.put(AWSConfigConstants.AWS_SECRET_ACCESS_KEY, "x");
  6. kinesisConfig.put(AWSConfigConstants.AWS_ENDPOINT, "http://localhost:4568");
  7. kinesisConfig.put("VerifyCertificate", "false");
  8. FlinkKinesisProducer<GenericRecord> kinesis = new FlinkKinesisProducer<GenericRecord>(new KinesisSerializationSchema<GenericRecord>() {
  9. @Override
  10. public ByteBuffer serialize(GenericRecord genericRecord) {
  11. return null;
  12. }
  13. @Override
  14. public String getTargetStream(GenericRecord genericRecord) {
  15. return null;
  16. }
  17. }, kinesisConfig);
  18. kinesis.setFailOnError(true);
  19. kinesis.setDefaultStream("myStreamName");
  20. kinesis.setDefaultPartition("0");

试着给当地的动情片做点什么,我一直在想

  1. Exception in thread "main" java.util.concurrent.ExecutionException: org.apache.flink.runtime.client.JobExecutionException: Job execution failed.
  2. at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)
  3. at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1895)
  4. at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1640)
  5. at org.apache.flink.streaming.api.environment.LocalStreamEnvironment.execute(LocalStreamEnvironment.java:74)
  6. at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1620)
  7. at org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1602)
  8. at ApacheFlinkPoc.main(ApacheFlinkPoc.java:80)
  9. Caused by: org.apache.flink.runtime.client.JobExecutionException: Job execution failed.
  10. at org.apache.flink.runtime.jobmaster.JobResult.toJobExecutionResult(JobResult.java:147)
  11. at org.apache.flink.client.program.PerJobMiniClusterFactory$PerJobMiniClusterJobClient.lambda$getJobExecutionResult$2(PerJobMiniClusterFactory.java:175)
  12. at java.util.concurrent.CompletableFuture.uniApply(CompletableFuture.java:602)
  13. at java.util.concurrent.CompletableFuture$UniApply.tryFire(CompletableFuture.java:577)
  14. at java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:474)
  15. at java.util.concurrent.CompletableFuture.complete(CompletableFuture.java:1962)
  16. at org.apache.flink.runtime.concurrent.FutureUtils$1.onComplete(FutureUtils.java:874)
  17. at akka.dispatch.OnComplete.internal(Future.scala:264)
  18. at akka.dispatch.OnComplete.internal(Future.scala:261)
  19. at akka.dispatch.japi$CallbackBridge.apply(Future.scala:191)
  20. at akka.dispatch.japi$CallbackBridge.apply(Future.scala:188)
  21. at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:60)
  22. at org.apache.flink.runtime.concurrent.Executors$DirectExecutionContext.execute(Executors.java:74)
  23. at scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:68)
  24. at scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1(Promise.scala:284)
  25. at scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1$adapted(Promise.scala:284)
  26. at scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:284)
  27. at akka.pattern.PromiseActorRef.$bang(AskSupport.scala:573)
  28. at akka.pattern.PipeToSupport$PipeableFuture$$anonfun$pipeTo$1.applyOrElse(PipeToSupport.scala:22)
  29. at akka.pattern.PipeToSupport$PipeableFuture$$anonfun$pipeTo$1.applyOrElse(PipeToSupport.scala:21)
  30. at scala.concurrent.Future.$anonfun$andThen$1(Future.scala:532)
  31. at scala.concurrent.impl.Promise.liftedTree1$1(Promise.scala:29)
  32. at scala.concurrent.impl.Promise.$anonfun$transform$1(Promise.scala:29)
  33. at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:60)
  34. at akka.dispatch.BatchingExecutor$AbstractBatch.processBatch(BatchingExecutor.scala:55)
  35. at akka.dispatch.BatchingExecutor$BlockableBatch.$anonfun$run$1(BatchingExecutor.scala:91)
  36. at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:12)
  37. at scala.concurrent.BlockContext$.withBlockContext(BlockContext.scala:81)
  38. at akka.dispatch.BatchingExecutor$BlockableBatch.run(BatchingExecutor.scala:91)
  39. at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:40)
  40. at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(ForkJoinExecutorConfigurator.scala:44)
  41. at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
  42. at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
  43. at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
  44. at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)
  45. Caused by: org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy
  46. at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:110)
  47. at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:76)
  48. at org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:192)
  49. at org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:186)
  50. at org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:180)
  51. at org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:496)
  52. at org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:380)
  53. at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
  54. at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
  55. at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
  56. at java.lang.reflect.Method.invoke(Method.java:498)
  57. at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:284)
  58. at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:199)
  59. at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)
  60. at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:152)
  61. at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26)
  62. at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21)
  63. at scala.PartialFunction.applyOrElse(PartialFunction.scala:123)
  64. at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122)
  65. at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)
  66. at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
  67. at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
  68. at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172)
  69. at akka.actor.Actor.aroundReceive(Actor.scala:517)
  70. at akka.actor.Actor.aroundReceive$(Actor.scala:515)
  71. at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225)
  72. at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592)
  73. at akka.actor.ActorCell.invoke(ActorCell.scala:561)
  74. at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258)
  75. at akka.dispatch.Mailbox.run(Mailbox.scala:225)
  76. at akka.dispatch.Mailbox.exec(Mailbox.scala:235)
  77. ... 4 more
  78. Caused by: java.lang.RuntimeException: An exception was thrown while processing a record: Unable to connect to endpoint
  79. Unable to connect to endpoint
  80. Unable to connect to endpoint
  81. Unable to connect to endpoint
  82. Record has reached expiration
  83. at org.apache.flink.streaming.connectors.kinesis.FlinkKinesisProducer.checkAndPropagateAsyncError(FlinkKinesisProducer.java:362)
  84. at org.apache.flink.streaming.connectors.kinesis.FlinkKinesisProducer.invoke(FlinkKinesisProducer.java:259)
  85. at org.apache.flink.streaming.api.operators.StreamSink.processElement(StreamSink.java:56)
  86. at org.apache.flink.streaming.runtime.tasks.OperatorChain$CopyingChainingOutput.pushToOperator(OperatorChain.java:641)
  87. at org.apache.flink.streaming.runtime.tasks.OperatorChain$CopyingChainingOutput.collect(OperatorChain.java:616)
  88. at org.apache.flink.streaming.runtime.tasks.OperatorChain$CopyingChainingOutput.collect(OperatorChain.java:596)
  89. at org.apache.flink.streaming.api.operators.AbstractStreamOperator$CountingOutput.collect(AbstractStreamOperator.java:730)
  90. at org.apache.flink.streaming.api.operators.AbstractStreamOperator$CountingOutput.collect(AbstractStreamOperator.java:708)
  91. at org.apache.flink.streaming.api.operators.StreamSourceContexts$NonTimestampContext.collect(StreamSourceContexts.java:104)
  92. at org.apache.flink.streaming.api.operators.StreamSourceContexts$NonTimestampContext.collectWithTimestamp(StreamSourceContexts.java:111)
  93. at org.apache.flink.streaming.connectors.kafka.internals.AbstractFetcher.emitRecordWithTimestamp(AbstractFetcher.java:398)
  94. at org.apache.flink.streaming.connectors.kafka.internal.KafkaFetcher.emitRecord(KafkaFetcher.java:185)
  95. at org.apache.flink.streaming.connectors.kafka.internal.KafkaFetcher.runFetchLoop(KafkaFetcher.java:150)
  96. at org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerBase.run(FlinkKafkaConsumerBase.java:718)
  97. at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:100)
  98. at org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:63)
  99. at org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:200)
  100. Caused by: org.apache.flink.kinesis.shaded.com.amazonaws.services.kinesis.producer.UserRecordFailedException
  101. at org.apache.flink.kinesis.shaded.com.amazonaws.services.kinesis.producer.KinesisProducer$MessageHandler.onPutRecordResult(KinesisProducer.java:197)
  102. at org.apache.flink.kinesis.shaded.com.amazonaws.services.kinesis.producer.KinesisProducer$MessageHandler.access$000(KinesisProducer.java:131)
  103. at org.apache.flink.kinesis.shaded.com.amazonaws.services.kinesis.producer.KinesisProducer$MessageHandler$1.run(KinesisProducer.java:138)
  104. at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
  105. at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
  106. at java.lang.Thread.run(Thread.java:748)

因为某种原因 flink-kinesis 连接器无法连接到 localhost:4568 有人能给你建议吗?
提前谢谢!!!

1hdlvixo

1hdlvixo1#

而不是通过-

  1. kinesisConfig.put(AWSConfigConstants.AWS_ENDPOINT, "http://localhost:4568");

像这样传球-

  1. kinesisConfig.put("KinesisEndpoint", "localhost")
  2. kinesisConfig.put("KinesisPort", "4567")
  3. kinesisConfig.put("VerifyCertificate","true")

相关问题