Commit b7e7b0e
committed
Make WindowedValue a public interface and receivers for it
When we introduce OutputBuilder, it will implement WindowedValue,
and every time we create an OutputBuilder we will want to provide
the following:
1. the default metadata for the output
2. the receiver for the output
This change unblocks 2 and reduces the size of the final change.1 parent 60307b4 commit b7e7b0e
File tree
474 files changed
+2622
-2643
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
- triggers
- 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
- adapter
- metrics
- translation
- functions
- types
- wrappers
- streaming
- io
- source
- bounded
- unbounded
- stableinput
- test/java/org/apache/beam/runners/flink
- streaming
- translation
- functions
- wrappers/streaming
- io
- source
- bounded
- unbounded
- stableinput
- google-cloud-dataflow-java
- src
- main/java/org/apache/beam/runners/dataflow
- util
- test/java/org/apache/beam/runners/dataflow
- util
- worker/src
- main/java/org/apache/beam/runners/dataflow/worker
- util
- windmill/work/processing
- test/java/org/apache/beam/runners/dataflow/worker
- graph
- util
- java-fn-execution/src
- main/java/org/apache/beam/runners/fnexecution
- control
- translation
- wire
- test/java/org/apache/beam/runners/fnexecution
- control
- data
- translation
- wire
- jet/src
- main/java/org/apache/beam/runners/jet
- processors
- test/java/org/apache/beam/runners/jet
- local-java/src/main/java/org/apache/beam/runners/local
- samza/src
- main/java/org/apache/beam/runners/samza
- adapter
- metrics
- runtime
- translation
- util
- test/java/org/apache/beam/runners/samza
- adapter
- metrics
- runtime
- spark
- 3/src
- main/java/org/apache/beam/runners/spark/structuredstreaming
- io
- translation
- batch
- functions
- helpers
- test/java/org/apache/beam/runners/spark/structuredstreaming/translation
- batch
- functions
- helpers
- src
- main/java/org/apache/beam/runners/spark
- coders
- io
- stateful
- translation
- streaming
- util
- test/java/org/apache/beam/runners/spark/translation
- twister2/src/main/java/org/apache/beam/runners/twister2
- translation/wrappers
- translators
- batch
- functions
- utils
- sdks/java
- core/src
- main/java/org/apache/beam/sdk
- transforms
- util
- construction
- values
- test/java/org/apache/beam/sdk
- fn/data
- transforms
- util
- construction
- extensions/euphoria/src/test/java/org/apache/beam/sdk/extensions/euphoria/core/testkit
- harness
- jmh/src/main/java/org/apache/beam/fn/harness/jmh
- src
- main/java/org/apache/beam/fn/harness
- control
- data
- debug
- state
- test/java/org/apache/beam/fn/harness
- data
- debug
- io
- components/src
- main/java/org/apache/beam/sdk/io/components/deadletterqueue
- sinks
- test/java/org/apache/beam/sdk/io/components/deadletterqueue
- google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/bigquery
- iceberg/src
- main/java/org/apache/beam/sdk/io/iceberg
- test/java/org/apache/beam/sdk/io/iceberg
Some content is hidden
Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
474 files changed
+2622
-2643
lines changedLines changed: 1 addition & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1233 | 1233 | | |
1234 | 1234 | | |
1235 | 1235 | | |
| 1236 | + | |
1236 | 1237 | | |
1237 | 1238 | | |
1238 | 1239 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
23 | | - | |
| 23 | + | |
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
| |||
Lines changed: 4 additions & 9 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
29 | 29 | | |
30 | 30 | | |
31 | 31 | | |
32 | | - | |
| 32 | + | |
33 | 33 | | |
34 | 34 | | |
35 | 35 | | |
| 36 | + | |
36 | 37 | | |
37 | 38 | | |
38 | 39 | | |
| |||
41 | 42 | | |
42 | 43 | | |
43 | 44 | | |
44 | | - | |
45 | | - | |
46 | | - | |
47 | | - | |
48 | | - | |
49 | | - | |
50 | 45 | | |
51 | 46 | | |
52 | 47 | | |
| |||
58 | 53 | | |
59 | 54 | | |
60 | 55 | | |
61 | | - | |
| 56 | + | |
62 | 57 | | |
63 | 58 | | |
64 | 59 | | |
| |||
168 | 163 | | |
169 | 164 | | |
170 | 165 | | |
171 | | - | |
| 166 | + | |
172 | 167 | | |
173 | 168 | | |
174 | 169 | | |
| |||
Lines changed: 5 additions & 31 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
17 | 17 | | |
18 | 18 | | |
19 | 19 | | |
20 | | - | |
21 | 20 | | |
22 | 21 | | |
23 | 22 | | |
24 | 23 | | |
25 | 24 | | |
26 | | - | |
27 | 25 | | |
28 | | - | |
| 26 | + | |
29 | 27 | | |
30 | 28 | | |
31 | 29 | | |
32 | 30 | | |
33 | | - | |
34 | 31 | | |
35 | 32 | | |
36 | 33 | | |
| |||
51 | 48 | | |
52 | 49 | | |
53 | 50 | | |
54 | | - | |
| 51 | + | |
55 | 52 | | |
56 | 53 | | |
57 | 54 | | |
| |||
68 | 65 | | |
69 | 66 | | |
70 | 67 | | |
71 | | - | |
| 68 | + | |
72 | 69 | | |
73 | 70 | | |
74 | 71 | | |
| |||
77 | 74 | | |
78 | 75 | | |
79 | 76 | | |
80 | | - | |
| 77 | + | |
81 | 78 | | |
82 | 79 | | |
83 | 80 | | |
| |||
91 | 88 | | |
92 | 89 | | |
93 | 90 | | |
94 | | - | |
95 | | - | |
96 | | - | |
97 | | - | |
98 | | - | |
99 | | - | |
100 | | - | |
101 | | - | |
102 | | - | |
103 | | - | |
104 | | - | |
105 | | - | |
106 | | - | |
107 | | - | |
108 | | - | |
109 | | - | |
110 | | - | |
111 | | - | |
112 | | - | |
113 | | - | |
114 | | - | |
115 | | - | |
116 | | - | |
117 | 91 | | |
118 | 92 | | |
119 | 93 | | |
| |||
130 | 104 | | |
131 | 105 | | |
132 | 106 | | |
133 | | - | |
| 107 | + | |
134 | 108 | | |
135 | 109 | | |
136 | 110 | | |
| |||
Lines changed: 2 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
29 | 29 | | |
30 | 30 | | |
31 | 31 | | |
32 | | - | |
33 | | - | |
34 | 32 | | |
35 | 33 | | |
| 34 | + | |
| 35 | + | |
36 | 36 | | |
37 | 37 | | |
38 | 38 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
18 | 18 | | |
19 | 19 | | |
20 | 20 | | |
21 | | - | |
| 21 | + | |
22 | 22 | | |
23 | 23 | | |
24 | 24 | | |
| |||
Lines changed: 2 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
28 | 28 | | |
29 | 29 | | |
30 | 30 | | |
31 | | - | |
32 | | - | |
| 31 | + | |
| 32 | + | |
33 | 33 | | |
34 | 34 | | |
35 | 35 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
20 | 20 | | |
21 | 21 | | |
22 | 22 | | |
23 | | - | |
| 23 | + | |
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
| |||
Lines changed: 3 additions & 2 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
25 | 25 | | |
26 | 26 | | |
27 | 27 | | |
28 | | - | |
29 | 28 | | |
| 29 | + | |
| 30 | + | |
30 | 31 | | |
31 | 32 | | |
32 | 33 | | |
| |||
140 | 141 | | |
141 | 142 | | |
142 | 143 | | |
143 | | - | |
| 144 | + | |
144 | 145 | | |
145 | 146 | | |
146 | 147 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
21 | 21 | | |
22 | 22 | | |
23 | 23 | | |
24 | | - | |
| 24 | + | |
25 | 25 | | |
26 | 26 | | |
27 | 27 | | |
| |||
0 commit comments