File tree
85 files changed
+1041
-643
lines changed- buildSrc/src/main/groovy/org/apache/beam/gradle
- runners
- core-java/src
- main/java/org/apache/beam/runners/core
- test/java/org/apache/beam/runners/core
- direct-java/src
- main/java/org/apache/beam/runners/direct
- test/java/org/apache/beam/runners/direct
- flink/src/main/java/org/apache/beam/runners/flink
- translation
- functions
- wrappers/streaming
- google-cloud-dataflow-java
- src/main/java/org/apache/beam/runners/dataflow
- worker/src
- main/java/org/apache/beam/runners/dataflow/worker
- util
- test/java/org/apache/beam/runners/dataflow/worker
- util
- java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control
- jet/src/main/java/org/apache/beam/runners/jet/processors
- samza/src/main/java/org/apache/beam/runners/samza/runtime
- spark
- 3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch
- functions
- src/main/java/org/apache/beam/runners/spark
- stateful
- translation
- twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions
- sdks/java
- core/src
- main/java/org/apache/beam/sdk
- testing
- transforms
- util
- construction
- values
- test/java/org/apache/beam/sdk/transforms
- reflect
- extensions/sql/zetasql/src/main/java/org/apache/beam/sdk/extensions/sql/zetasql
- harness/src
- main/java/org/apache/beam/fn/harness
- test/java/org/apache/beam/fn/harness
- io
- google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp
- bigquery
- spanner
- kafka/src/test/java/org/apache/beam/sdk/io/kafka
- pulsar/src/test/java/org/apache/beam/sdk/io/pulsar
- sparkreceiver/3/src/test/java/org/apache/beam/sdk/io/sparkreceiver
Some content is hidden
Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
85 files changed
+1041
-643
lines changedLines changed: 3 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1198 | 1198 | | |
1199 | 1199 | | |
1200 | 1200 | | |
| 1201 | + | |
1201 | 1202 | | |
1202 | 1203 | | |
1203 | 1204 | | |
| |||
1291 | 1292 | | |
1292 | 1293 | | |
1293 | 1294 | | |
1294 | | - | |
| 1295 | + | |
| 1296 | + | |
1295 | 1297 | | |
1296 | 1298 | | |
1297 | 1299 | | |
| |||
Lines changed: 3 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
25 | 25 | | |
26 | 26 | | |
27 | 27 | | |
| 28 | + | |
28 | 29 | | |
29 | 30 | | |
30 | 31 | | |
| |||
91 | 92 | | |
92 | 93 | | |
93 | 94 | | |
94 | | - | |
95 | | - | |
| 95 | + | |
| 96 | + | |
96 | 97 | | |
97 | 98 | | |
98 | 99 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
141 | 141 | | |
142 | 142 | | |
143 | 143 | | |
144 | | - | |
| 144 | + | |
145 | 145 | | |
146 | 146 | | |
147 | 147 | | |
| |||
Lines changed: 7 additions & 6 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
45 | 45 | | |
46 | 46 | | |
47 | 47 | | |
| 48 | + | |
48 | 49 | | |
49 | 50 | | |
50 | 51 | | |
| |||
72 | 73 | | |
73 | 74 | | |
74 | 75 | | |
75 | | - | |
| 76 | + | |
76 | 77 | | |
77 | 78 | | |
78 | 79 | | |
| |||
98 | 99 | | |
99 | 100 | | |
100 | 101 | | |
101 | | - | |
| 102 | + | |
102 | 103 | | |
103 | 104 | | |
104 | 105 | | |
| |||
375 | 376 | | |
376 | 377 | | |
377 | 378 | | |
378 | | - | |
| 379 | + | |
379 | 380 | | |
380 | 381 | | |
381 | 382 | | |
| |||
390 | 391 | | |
391 | 392 | | |
392 | 393 | | |
393 | | - | |
| 394 | + | |
394 | 395 | | |
395 | 396 | | |
396 | 397 | | |
| |||
403 | 404 | | |
404 | 405 | | |
405 | 406 | | |
406 | | - | |
| 407 | + | |
407 | 408 | | |
408 | 409 | | |
409 | 410 | | |
| |||
413 | 414 | | |
414 | 415 | | |
415 | 416 | | |
416 | | - | |
| 417 | + | |
417 | 418 | | |
418 | 419 | | |
419 | 420 | | |
| |||
Lines changed: 3 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
46 | 46 | | |
47 | 47 | | |
48 | 48 | | |
| 49 | + | |
49 | 50 | | |
50 | 51 | | |
51 | 52 | | |
| |||
106 | 107 | | |
107 | 108 | | |
108 | 109 | | |
109 | | - | |
| 110 | + | |
110 | 111 | | |
111 | 112 | | |
112 | 113 | | |
| |||
214 | 215 | | |
215 | 216 | | |
216 | 217 | | |
217 | | - | |
| 218 | + | |
218 | 219 | | |
219 | 220 | | |
220 | 221 | | |
| |||
Lines changed: 12 additions & 4 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
404 | 404 | | |
405 | 405 | | |
406 | 406 | | |
407 | | - | |
| 407 | + | |
408 | 408 | | |
409 | 409 | | |
410 | 410 | | |
| |||
436 | 436 | | |
437 | 437 | | |
438 | 438 | | |
439 | | - | |
| 439 | + | |
440 | 440 | | |
441 | 441 | | |
442 | 442 | | |
| |||
446 | 446 | | |
447 | 447 | | |
448 | 448 | | |
449 | | - | |
450 | | - | |
| 449 | + | |
| 450 | + | |
| 451 | + | |
| 452 | + | |
| 453 | + | |
| 454 | + | |
| 455 | + | |
| 456 | + | |
| 457 | + | |
| 458 | + | |
451 | 459 | | |
452 | 460 | | |
453 | 461 | | |
| |||
runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java
Lines changed: 3 additions & 3 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
430 | 430 | | |
431 | 431 | | |
432 | 432 | | |
433 | | - | |
| 433 | + | |
434 | 434 | | |
435 | 435 | | |
436 | 436 | | |
| |||
490 | 490 | | |
491 | 491 | | |
492 | 492 | | |
493 | | - | |
| 493 | + | |
494 | 494 | | |
495 | 495 | | |
496 | 496 | | |
| |||
544 | 544 | | |
545 | 545 | | |
546 | 546 | | |
547 | | - | |
| 547 | + | |
548 | 548 | | |
549 | 549 | | |
550 | 550 | | |
| |||
Lines changed: 2 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
41 | 41 | | |
42 | 42 | | |
43 | 43 | | |
| 44 | + | |
44 | 45 | | |
45 | 46 | | |
46 | 47 | | |
| |||
112 | 113 | | |
113 | 114 | | |
114 | 115 | | |
115 | | - | |
| 116 | + | |
116 | 117 | | |
117 | 118 | | |
118 | 119 | | |
| |||
Lines changed: 29 additions & 21 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
759 | 759 | | |
760 | 760 | | |
761 | 761 | | |
762 | | - | |
| 762 | + | |
| 763 | + | |
763 | 764 | | |
764 | 765 | | |
765 | 766 | | |
| |||
805 | 806 | | |
806 | 807 | | |
807 | 808 | | |
808 | | - | |
| 809 | + | |
| 810 | + | |
809 | 811 | | |
810 | 812 | | |
811 | 813 | | |
| |||
846 | 848 | | |
847 | 849 | | |
848 | 850 | | |
849 | | - | |
| 851 | + | |
| 852 | + | |
850 | 853 | | |
851 | 854 | | |
852 | 855 | | |
| |||
879 | 882 | | |
880 | 883 | | |
881 | 884 | | |
882 | | - | |
| 885 | + | |
883 | 886 | | |
884 | 887 | | |
885 | 888 | | |
| |||
1486 | 1489 | | |
1487 | 1490 | | |
1488 | 1491 | | |
1489 | | - | |
| 1492 | + | |
| 1493 | + | |
1490 | 1494 | | |
1491 | 1495 | | |
1492 | 1496 | | |
| |||
1527 | 1531 | | |
1528 | 1532 | | |
1529 | 1533 | | |
1530 | | - | |
| 1534 | + | |
| 1535 | + | |
1531 | 1536 | | |
1532 | 1537 | | |
1533 | 1538 | | |
| |||
1614 | 1619 | | |
1615 | 1620 | | |
1616 | 1621 | | |
1617 | | - | |
| 1622 | + | |
| 1623 | + | |
1618 | 1624 | | |
1619 | 1625 | | |
1620 | 1626 | | |
| |||
1661 | 1667 | | |
1662 | 1668 | | |
1663 | 1669 | | |
1664 | | - | |
| 1670 | + | |
1665 | 1671 | | |
1666 | 1672 | | |
1667 | 1673 | | |
| |||
1702 | 1708 | | |
1703 | 1709 | | |
1704 | 1710 | | |
1705 | | - | |
| 1711 | + | |
1706 | 1712 | | |
1707 | 1713 | | |
1708 | 1714 | | |
| |||
1711 | 1717 | | |
1712 | 1718 | | |
1713 | 1719 | | |
1714 | | - | |
| 1720 | + | |
| 1721 | + | |
1715 | 1722 | | |
1716 | 1723 | | |
1717 | 1724 | | |
| |||
1811 | 1818 | | |
1812 | 1819 | | |
1813 | 1820 | | |
1814 | | - | |
| 1821 | + | |
1815 | 1822 | | |
1816 | 1823 | | |
1817 | 1824 | | |
| |||
1850 | 1857 | | |
1851 | 1858 | | |
1852 | 1859 | | |
1853 | | - | |
| 1860 | + | |
| 1861 | + | |
1854 | 1862 | | |
1855 | 1863 | | |
1856 | 1864 | | |
| |||
1868 | 1876 | | |
1869 | 1877 | | |
1870 | 1878 | | |
1871 | | - | |
| 1879 | + | |
1872 | 1880 | | |
1873 | 1881 | | |
1874 | 1882 | | |
| |||
2193 | 2201 | | |
2194 | 2202 | | |
2195 | 2203 | | |
2196 | | - | |
2197 | | - | |
| 2204 | + | |
| 2205 | + | |
2198 | 2206 | | |
2199 | 2207 | | |
2200 | 2208 | | |
2201 | 2209 | | |
2202 | 2210 | | |
2203 | 2211 | | |
2204 | 2212 | | |
2205 | | - | |
2206 | | - | |
| 2213 | + | |
| 2214 | + | |
2207 | 2215 | | |
2208 | 2216 | | |
2209 | 2217 | | |
| |||
2231 | 2239 | | |
2232 | 2240 | | |
2233 | 2241 | | |
2234 | | - | |
2235 | | - | |
| 2242 | + | |
| 2243 | + | |
2236 | 2244 | | |
2237 | 2245 | | |
2238 | 2246 | | |
2239 | 2247 | | |
2240 | 2248 | | |
2241 | 2249 | | |
2242 | 2250 | | |
2243 | | - | |
2244 | | - | |
| 2251 | + | |
| 2252 | + | |
2245 | 2253 | | |
2246 | 2254 | | |
2247 | 2255 | | |
| |||
0 commit comments