apache / apache/beam

Streaming test hangs

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

Description

More information on SO:
[https://stackoverflow.com/questions/49266481/unit-test-hangs-forever-if-dofn-resets-event-timers](https://stackoverflow.com/questions/49266481/unit-test-hangs-forever-if-dofn-resets-event-timers)

Test hangs when run despite fairly simple semantics. Modified a wordcount test with the code in the SO to generate a full example.

 
```

package com.example;

import org.apache.beam.sdk.testing.TestStream;
import java.util.Arrays;
import
java.util.List;
import org.joda.time.Instant;
import org.joda.time.Duration;
import org.apache.beam.sdk.transforms.ParDo;
import
org.apache.beam.sdk.values.TimestampedValue;
import com.example.WordCount.CountWords;
import com.example.WordCount.ExtractWordsFn;
import
com.example.WordCount.FormatAsTextFn;
import org.apache.beam.sdk.state.Timer;
import org.apache.beam.sdk.state.TimerSpecs;
import
org.apache.beam.sdk.state.TimeDomain;
import org.apache.beam.sdk.state.TimerSpec;
import org.apache.beam.sdk.values.KV;
import
org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.testing.PAssert;
import org.apache.beam.sdk.testing.TestPipeline;
import
org.apache.beam.sdk.testing.ValidatesRunner;
import org.apache.beam.sdk.transforms.Create;
import
org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.DoFnTester;
import org.apache.beam.sdk.transforms.MapElements;
import
org.apache.beam.sdk.values.PCollection;
import org.hamcrest.CoreMatchers;
import org.junit.Assert;
import
org.junit.Rule;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import
org.junit.runner.RunWith;
import org.junit.runners.JUnit4;

/**
* Tests of WordCount.
*/
@RunWith(JUnit4.class)
public
class WordCountTest {

static class KeyElements extends DoFn> {
   @ProcessElement
   public
void processElement(ProcessContext context) {
       final String[] parts = context.element().split(":");
       if
(parts.length == 2) {
           context.output(KV.of(parts[0], parts[1]));
       }
   }
}

static
class TimerDoFn extends DoFn, KV> {
   @TimerId("expiry")
   private
final TimerSpec timerSpec = TimerSpecs.timer(TimeDomain.EVENT_TIME);

   @ProcessElement
   public
void processElement(ProcessContext context, @TimerId("expiry") Timer timer) {
       timer.set(context.timestamp().plus(Duration.standardHours(1)));
       final
KV e = context.element();
       context.output(KV.of(e.getKey(), e.getValue() + "_output"));
   }

   @OnTimer("expiry")
   public
void onExpiry(OnTimerContext context) {
       // do nothing
   }
}

 @Rule
 public TestPipeline
p = TestPipeline.create();

 /** Example test that tests a PTransform by using an in-memory input
and inspecting the output. */
 @Test
 @Category(ValidatesRunner.class)
 public void testCountWords()
throws Exception {
TestStream stream = TestStream
           .create(StringUtf8Coder.of())
           .addElements(
                   TimestampedValue.of("a:0",
new Instant(0)),
                   TimestampedValue.of("a:1", new Instant(1)),
                   TimestampedValue.of("a:2",
new Instant(2)),
                   TimestampedValue.of("a:3", new Instant(3)))
           .advanceWatermarkToInfinity();

   PCollection> result = p
           .apply(stream)
           .apply(ParDo.of(new KeyElements()))
           .apply(ParDo.of(new
TimerDoFn()));

   PAssert.that(result).containsInAnyOrder(
           KV.of("a", "0_output"),
           KV.of("a",
"1_output"),
           KV.of("a", "2_output"),
           KV.of("a", "3_output"));

   p.run();
 
}
}

```

Imported from Jira [BEAM-3880](https://issues.apache.org/jira/browse/BEAM-3880). Original Jira may contain additional context.
Reported by: larasch.

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.