Hi Ankit,

I guess it can be a bug there. Let me check this out.

> On 20 Sep 2019, at 00:12, Ankit Jhalaria <[email protected]> wrote:
> 
> Hey beam devs,
> 
> I am using beam 2.15 and while doing KinesisIO.write() getting a NPE. 
> This is how I am using it:
> KinesisIO.write()
>     .withStreamName(“streamName")
>     .withPartitionKey("DEFAULT_PARTITION")
>     .withAWSClientsProvider(
>         “A”, “B”, Regions.US_WEST_2)
> 
> More specifically 
> https://github.com/apache/beam/blob/master/sdks/java/io/kinesis/src/main/java/org/apache/beam/sdk/io/kinesis/KinesisIO.java#L685
>  
> <https://github.com/apache/beam/blob/master/sdks/java/io/kinesis/src/main/java/org/apache/beam/sdk/io/kinesis/KinesisIO.java#L685>
>  is failing with NPE
> 
> I do see that the producer is transient
> private static transient IKinesisProducer producer;
> I was wondering if that had do anything with the error below.
> 
> org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapper.run(UnboundedSourceWrapper.java:297)\n\tat
>  
> org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:93)\n\tat
>  
> org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:57)\n\tat
>  
> org.apache.flink.streaming.runtime.tasks.SourceStreamTask.run(SourceStreamTask.java:97)\n\tat
>  
> org.apache.flink.streaming.runtime.tasks.StoppableSourceStreamTask.run(StoppableSourceStreamTask.java:45)\n\tat
>  
> org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:300)\n\tat
>  org.apache.flink.runtime.taskmanager.Task.run(Task.java:711)\n\tat 
> java.lang.Thread.run(Thread.java:748)\nCaused by: 
> java.lang.NullPointerException: null\n\tat 
> org.apache.beam.sdk.io.kinesis.KinesisIO$Write$KinesisWriterFn.flushBundle(KinesisIO.java:685)\n\tat
>  
> org.apache.beam.sdk.io.kinesis.KinesisIO$Write$KinesisWriterFn.finishBundle(KinesisIO.java:669)
> 
> Best,
> Ankit

Reply via email to