apache / apache/beam

Direct Runner State is null while active timers exist

Open
#20,966 0 comments 0 reactions 0 assignees View on GitHub
bug direct P3 runners
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

State is set to `null` while active timer is present, this issue does not show in other runners.

The following example will reach the IllegalStateException within 10-20 times of it being run. `LOOP_COUNT` does not seem to be a factor as it reproduces with 100 or 100000 `LOOP_COUNT`. The number of keys is a factor as it did not reproduce with only one key, have not tried with more than 3 keys to see if it's easier to reproduce. 
 
```

package test;

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.coders.BigEndianIntegerCoder;
import
org.apache.beam.sdk.coders.KvCoder;
import org.apache.beam.sdk.state.StateSpec;
import org.apache.beam.sdk.state.StateSpecs;
import
org.apache.beam.sdk.state.TimeDomain;
import org.apache.beam.sdk.state.Timer;
import org.apache.beam.sdk.state.TimerSpec;
import
org.apache.beam.sdk.state.TimerSpecs;
import org.apache.beam.sdk.state.ValueState;
import org.apache.beam.sdk.testing.TestStream;
import
org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
import
org.apache.beam.sdk.transforms.WithKeys;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
import
org.joda.time.Duration;
import org.joda.time.Instant;

import java.util.Optional;
 

public class
Test {

   public static void main (String [] args) throws Exception{
       Test.testToFailure();

   }

   public
static void testToFailure() throws Exception {
       int count = 0;

       while (true) {
           failingTest();
           System.out.println(
                   String.format("Got
to Count %s", String.valueOf(count++)));
       }
   }

   public static void failingTest() throws
Exception {
       Pipeline p = Pipeline.create();

       Instant now = Instant.now();
       TestStream
stream =
               TestStream.create(BigEndianIntegerCoder.of())
                       .addElements(1)
                       .advanceWatermarkTo(now.plus(Duration.standardSeconds(1)))
                       .addElements(2)
                       .advanceWatermarkTo(now.plus(Duration.standardSeconds(1)))
                       .addElements(3)
                       .advanceWatermarkToInfinity();

       p.apply(stream)
               .apply(WithKeys.of(x
-> x))
               .setCoder(KvCoder.of(BigEndianIntegerCoder.of(), BigEndianIntegerCoder.of()))
               .apply(new
TestToFail());
       p.run();
   }

   public static class TestToFail
           extends PTransform>, PCollection> {

       @Override
       public PCollection expand(PCollection> input) {
           return input.apply(ParDo.of(new LoopingRead()));
       }
   }

   public
static class LoopingRead extends DoFn, Integer> {

       static int LOOP_COUNT
= 100;

       @StateId("value")
       private final StateSpec> value =
               StateSpecs.value(BigEndianIntegerCoder.of());

       @StateId("count")
       private
final StateSpec> count =
               StateSpecs.value(BigEndianIntegerCoder.of());

       @TimerId("actionTimers")
       private
final TimerSpec timer = TimerSpecs.timer(TimeDomain.EVENT_TIME);

       @ProcessElement
       public
void processElement(
               ProcessContext c,
               @StateId("value") ValueState
value,
               @TimerId("actionTimers") Timer timers) {

           value.write(c.element().getValue());
           timers.set(c.timestamp().plus(Duration.millis(1000)));
       }

       /**
*/
       @OnTimer("actionTimers")
       public void onTimer(
               OnTimerContext c,
               @StateId("value")
ValueState value,
               @StateId("count") ValueState count,
               @TimerId("actionTimers")
Timer timers) {

           if (value.read() == null) {
               throw new IllegalStateException("BINGO!");
           }
           Integer
counter = Optional.ofNullable(count.read()).orElse(0) + 1;
           count.write(counter);
           value.write(value.read()
+ counter);

           if (counter < LOOP_COUNT) {
               timers.set(c.timestamp().plus(Duration.standardSeconds(1)));
           }
       }
   }
}

```

Imported from Jira [BEAM-11971](https://issues.apache.org/jira/browse/BEAM-11971). Original Jira may contain additional context.
Reported by: rarokni@gmail.com.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.