[ 
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)

Reply via email to