[
https://issues.apache.org/jira/browse/BEAM-1229?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=15818288#comment-15818288
]
ASF GitHub Bot commented on BEAM-1229:
--------------------------------------
GitHub user xhumanoid opened a pull request:
https://github.com/apache/beam/pull/1765
BEAM-1229 flink KafkaIOExamples submit error
Behavior changed when TypedPValue started validating input coders
038950d 22.12.16 Add Parameters to finishSpecifying
setup coder for input in correct value
You can merge this pull request into a Git repository by running:
$ git pull https://github.com/xhumanoid/incubator-beam
BEAM-1229-flink-KafkaIOExamples-error
Alternatively you can review and apply these changes as the patch at:
https://github.com/apache/beam/pull/1765.patch
To close this pull request, make a commit to your master/trunk branch
with (at least) the following in the commit message:
This closes #1765
----
commit 4aa9466e4dc5e7db8ee92fe7f2aa7fedfe41f42a
Author: Alexey Diomin <[email protected]>
Date: 2017-01-11T13:08:35Z
BEAM-1229 flink KafkaIOExamples submit error
----
> flink KafkaIOExamples submit error
> -----------------------------------
>
> Key: BEAM-1229
> URL: https://issues.apache.org/jira/browse/BEAM-1229
> Project: Beam
> Issue Type: Bug
> Components: runner-flink
> Affects Versions: 0.5.0
> Environment: Flink:1.1.3_2.11
> JDK:Oracle jdk 1.8.0_73
> Beam:0.5.0-incubator-SNAPSHOT
> OS:Windows 10
> Reporter: Fei Feng
> Assignee: Maximilian Michels
> Priority: Minor
> Labels: Flink, FlinkRunner, KafkaIO
>
> I change all the beam pom.xml scala to 2.11.8,scala lib to 2.11,compile the
> beam jars。
> Submit KafkaIOExamples in runners/FlinkRunner/KafkaIOExamples to flink,main
> class "ReadStringFromKafka"
> error occured as follow:
> org.apache.flink.client.program.ProgramInvocationException: The main method
> caused an error.
> at
> org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:524)
> at
> org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:403)
> at
> org.apache.flink.client.program.OptimizerPlanEnvironment.getOptimizedPlan(OptimizerPlanEnvironment.java:80)
> at
> org.apache.flink.client.program.ClusterClient.getOptimizedPlan(ClusterClient.java:274)
> at
> org.apache.flink.runtime.webmonitor.handlers.JarActionHandler.getJobGraphAndClassLoader(JarActionHandler.java:95)
> at
> org.apache.flink.runtime.webmonitor.handlers.JarPlanHandler.handleRequest(JarPlanHandler.java:42)
> at
> org.apache.flink.runtime.webmonitor.RuntimeMonitorHandler.respondAsLeader(RuntimeMonitorHandler.java:88)
> at
> org.apache.flink.runtime.webmonitor.RuntimeMonitorHandlerBase.channelRead0(RuntimeMonitorHandlerBase.java:84)
> at
> org.apache.flink.runtime.webmonitor.RuntimeMonitorHandlerBase.channelRead0(RuntimeMonitorHandlerBase.java:44)
> at
> io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
> at
> io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:339)
> at
> io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:324)
> at io.netty.handler.codec.http.router.Handler.routed(Handler.java:62)
> at
> io.netty.handler.codec.http.router.DualAbstractHandler.channelRead0(DualAbstractHandler.java:57)
> at
> io.netty.handler.codec.http.router.DualAbstractHandler.channelRead0(DualAbstractHandler.java:20)
> at
> io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
> at
> io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:339)
> at
> io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:324)
> at
> org.apache.flink.runtime.webmonitor.HttpRequestHandler.channelRead0(HttpRequestHandler.java:105)
> at
> org.apache.flink.runtime.webmonitor.HttpRequestHandler.channelRead0(HttpRequestHandler.java:65)
> at
> io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
> at
> io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:339)
> at
> io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:324)
> at
> io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:242)
> at
> io.netty.channel.CombinedChannelDuplexHandler.channelRead(CombinedChannelDuplexHandler.java:147)
> at
> io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:339)
> at
> io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:324)
> at
> io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:847)
> at
> io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:131)
> at
> io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:511)
> at
> io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:468)
> at
> io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:382)
> at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:354)
> at
> io.netty.util.concurrent.SingleThreadEventExecutor$2.run(SingleThreadEventExecutor.java:111)
> at
> io.netty.util.concurrent.DefaultThreadFactory$DefaultRunnableDecorator.run(DefaultThreadFactory.java:137)
> at java.lang.Thread.run(Thread.java:745)
> Caused by: java.lang.IllegalStateException: null
> at
> org.apache.beam.sdk.repackaged.com.google.common.base.Preconditions.checkState(Preconditions.java:174)
> at org.apache.beam.sdk.values.TypedPValue.getCoder(TypedPValue.java:51)
> at org.apache.beam.sdk.values.PCollection.getCoder(PCollection.java:130)
> at
> org.apache.beam.sdk.values.TypedPValue.finishSpecifying(TypedPValue.java:90)
> at
> org.apache.beam.sdk.runners.TransformHierarchy.finishSpecifyingInput(TransformHierarchy.java:95)
> at org.apache.beam.sdk.Pipeline.applyInternal(Pipeline.java:386)
> at org.apache.beam.sdk.Pipeline.applyTransform(Pipeline.java:302)
> at org.apache.beam.sdk.values.PCollection.apply(PCollection.java:154)
> at
> com.xdata.beam.wordcount.KafkaIOExamples$KafkaString$ReadStringFromKafka.main(KafkaIOExamples.java:83)
> 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:497)
> at
> org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:509)
> ... 35 more
--
This message was sent by Atlassian JIRA
(v6.3.4#6332)