Streaming test hangs
- 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
Assessment
This issue has not been assessed yet.