kennknowles opened a new issue, #18709:
URL: https://github.com/apache/beam/issues/18709
assert_that does not work for AfterWatermark timers.
Easy way to reproduce: modify test_gbk_execution [1] in this form:
```
def test_this(self):
test_stream = (TestStream()
.add_elements(['a', 'b',
'c'])
.advance_watermark_to(20))
def fnc(x):
print 'fired_elem:',
x
return x
options = PipelineOptions()
options.view_as(StandardOptions).streaming
= True
p = TestPipeline(options=options)
records = (p
| test_stream
| beam.WindowInto(
FixedWindows(15),
trigger=trigger.AfterWatermark(early=trigger.AfterCount(2)),
accumulation_mode=trigger.AccumulationMode.ACCUMULATING)
| beam.Map(lambda
x: ('k', x))
| beam.GroupByKey())
assert_that(records, equal_to([
('k',
['a', 'b', 'c'])]))
p.run()
```
This test will pass, but if the .advance_watermark_to(20) is removed, the
test will fail. However, both cases fire the same elements:
fired_elem: ('k', ['a', 'b', 'c'])
fired_elem: ('k', ['a', 'b', 'c'])
In the passing case, they correspond to the sorted_actual inside the
assert_that. In the failing case:
sorted_actual: [('k', ['a', 'b', 'c']), ('k', ['a', 'b', 'c'])]
sorted_actual: []
[1]
https://github.com/mariapython/incubator-beam/blob/direct-timers-show/sdks/python/apache_beam/testing/test_stream_test.py#L120
Imported from Jira
[BEAM-3377](https://issues.apache.org/jira/browse/BEAM-3377). Original Jira may
contain additional context.
Reported by: mariagh.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]