Commit 1ae3fe7
committed
[Python] Add UnboundedSource ValidatesRunner test on portable runners
Add test_unbounded_source_read to PortableRunnerTest so the portable
ValidatesRunner suites exercise the UnboundedSource SDF wrapper end to
end: read a self-terminating source through the job service, assert the
elements and that the EOF MAX_TIMESTAMP watermark lets a downstream
FixedWindows + GroupByKey fire.
The embedded portable runner variants, the Flink suites, and
PrismRunnerTest inherit the test. SparkRunnerTest skips it because
portable Spark does not execute SDFs (#19468), matching its other SDF
test skips.1 parent 6e6a40c commit 1ae3fe7
1 file changed
Lines changed: 26 additions & 0 deletions
Lines changed: 26 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
30 | 30 | | |
31 | 31 | | |
32 | 32 | | |
| 33 | + | |
33 | 34 | | |
34 | 35 | | |
35 | 36 | | |
| |||
47 | 48 | | |
48 | 49 | | |
49 | 50 | | |
| 51 | + | |
50 | 52 | | |
51 | 53 | | |
52 | 54 | | |
| |||
209 | 211 | | |
210 | 212 | | |
211 | 213 | | |
| 214 | + | |
| 215 | + | |
| 216 | + | |
| 217 | + | |
| 218 | + | |
| 219 | + | |
| 220 | + | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
| 235 | + | |
| 236 | + | |
| 237 | + | |
212 | 238 | | |
213 | 239 | | |
214 | 240 | | |
| |||
0 commit comments