Skip to content

Commit 91d3b6b

Browse files
committed
move enum to top level class
1 parent 4e76198 commit 91d3b6b

28 files changed

Lines changed: 152 additions & 194 deletions
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
package org.apache.beam.runners.core;
19+
20+
public enum CausedByDrain {
21+
CAUSED_BY_DRAIN,
22+
NORMAL
23+
}

runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnContextFactory.java

Lines changed: 3 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -136,22 +136,19 @@ public TimersImpl(StateNamespace namespace) {
136136
@Override
137137
public void setTimer(Instant timestamp, TimeDomain timeDomain) {
138138
timerInternals.setTimer(
139-
TimerData.of(
140-
namespace, timestamp, timestamp, timeDomain, TimerData.CausedByDrain.NORMAL));
139+
TimerData.of(namespace, timestamp, timestamp, timeDomain, CausedByDrain.NORMAL));
141140
}
142141

143142
@Override
144143
public void setTimer(Instant timestamp, Instant outputTimestamp, TimeDomain timeDomain) {
145144
timerInternals.setTimer(
146-
TimerData.of(
147-
namespace, timestamp, outputTimestamp, timeDomain, TimerData.CausedByDrain.NORMAL));
145+
TimerData.of(namespace, timestamp, outputTimestamp, timeDomain, CausedByDrain.NORMAL));
148146
}
149147

150148
@Override
151149
public void deleteTimer(Instant timestamp, TimeDomain timeDomain) {
152150
timerInternals.deleteTimer(
153-
TimerData.of(
154-
namespace, timestamp, timestamp, timeDomain, TimerData.CausedByDrain.NORMAL));
151+
TimerData.of(namespace, timestamp, timestamp, timeDomain, CausedByDrain.NORMAL));
155152
}
156153

157154
@Override

runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -604,7 +604,7 @@ public String getErrorContext() {
604604
wakeupTime,
605605
wakeupTime,
606606
TimeDomain.PROCESSING_TIME,
607-
TimerInternals.TimerData.CausedByDrain.NORMAL));
607+
CausedByDrain.NORMAL));
608608
}
609609

610610
private DoFnInvoker.ArgumentProvider<InputT, OutputT> wrapOptionsAsSetup(

runners/core-java/src/main/java/org/apache/beam/runners/core/TimerInternals.java

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -168,10 +168,6 @@ void deleteTimer(
168168
/** Data about a timer as represented within {@link TimerInternals}. */
169169
@AutoValue
170170
abstract class TimerData implements Comparable<TimerData> {
171-
public enum CausedByDrain {
172-
CAUSED_BY_DRAIN,
173-
NORMAL
174-
}
175171

176172
public abstract String getTimerId();
177173

@@ -245,7 +241,7 @@ public static TimerData of(
245241
timestamp,
246242
outputTimestamp,
247243
domain,
248-
TimerData.CausedByDrain.NORMAL);
244+
CausedByDrain.NORMAL);
249245
}
250246

251247
/**
@@ -355,7 +351,7 @@ public TimerData decode(InputStream inStream) throws CoderException, IOException
355351
timestamp,
356352
outputTimestamp,
357353
domain,
358-
TimerData.CausedByDrain.NORMAL);
354+
CausedByDrain.NORMAL);
359355
}
360356

361357
@Override
@@ -401,8 +397,7 @@ public TimerData decode(InputStream inStream) throws CoderException, IOException
401397
StateNamespaces.fromString(STRING_CODER.decode(inStream), windowCoder);
402398
Instant timestamp = INSTANT_CODER.decode(inStream);
403399
TimeDomain domain = TimeDomain.valueOf(STRING_CODER.decode(inStream));
404-
return TimerData.of(
405-
timerId, namespace, timestamp, timestamp, domain, TimerData.CausedByDrain.NORMAL);
400+
return TimerData.of(timerId, namespace, timestamp, timestamp, domain, CausedByDrain.NORMAL);
406401
}
407402

408403
@Override

runners/core-java/src/test/java/org/apache/beam/runners/core/InMemoryTimerInternalsTest.java

Lines changed: 12 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -47,15 +47,15 @@ public void testFiringEventTimers() throws Exception {
4747
new Instant(19),
4848
new Instant(19),
4949
TimeDomain.EVENT_TIME,
50-
TimerData.CausedByDrain.NORMAL);
50+
CausedByDrain.NORMAL);
5151
TimerData eventTimer2 =
5252
TimerData.of(
5353
ID2,
5454
NS1,
5555
new Instant(29),
5656
new Instant(29),
5757
TimeDomain.EVENT_TIME,
58-
TimerData.CausedByDrain.NORMAL);
58+
CausedByDrain.NORMAL);
5959

6060
underTest.setTimer(eventTimer1);
6161
underTest.setTimer(eventTimer2);
@@ -128,14 +128,14 @@ public void testFiringProcessingTimeTimers() throws Exception {
128128
new Instant(19),
129129
new Instant(19),
130130
TimeDomain.PROCESSING_TIME,
131-
TimerData.CausedByDrain.NORMAL);
131+
CausedByDrain.NORMAL);
132132
TimerData processingTime2 =
133133
TimerData.of(
134134
NS1,
135135
new Instant(29),
136136
new Instant(29),
137137
TimeDomain.PROCESSING_TIME,
138-
TimerData.CausedByDrain.NORMAL);
138+
CausedByDrain.NORMAL);
139139

140140
underTest.setTimer(processingTime1);
141141
underTest.setTimer(processingTime2);
@@ -165,46 +165,38 @@ public void testTimerOrdering() throws Exception {
165165
InMemoryTimerInternals underTest = new InMemoryTimerInternals();
166166
TimerData eventTime1 =
167167
TimerData.of(
168-
NS1,
169-
new Instant(19),
170-
new Instant(19),
171-
TimeDomain.EVENT_TIME,
172-
TimerData.CausedByDrain.NORMAL);
168+
NS1, new Instant(19), new Instant(19), TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
173169
TimerData processingTime1 =
174170
TimerData.of(
175171
NS1,
176172
new Instant(19),
177173
new Instant(19),
178174
TimeDomain.PROCESSING_TIME,
179-
TimerData.CausedByDrain.NORMAL);
175+
CausedByDrain.NORMAL);
180176
TimerData synchronizedProcessingTime1 =
181177
TimerData.of(
182178
NS1,
183179
new Instant(19),
184180
new Instant(19),
185181
TimeDomain.SYNCHRONIZED_PROCESSING_TIME,
186-
TimerData.CausedByDrain.NORMAL);
182+
CausedByDrain.NORMAL);
187183
TimerData eventTime2 =
188184
TimerData.of(
189-
NS1,
190-
new Instant(29),
191-
new Instant(29),
192-
TimeDomain.EVENT_TIME,
193-
TimerData.CausedByDrain.NORMAL);
185+
NS1, new Instant(29), new Instant(29), TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
194186
TimerData processingTime2 =
195187
TimerData.of(
196188
NS1,
197189
new Instant(29),
198190
new Instant(29),
199191
TimeDomain.PROCESSING_TIME,
200-
TimerData.CausedByDrain.NORMAL);
192+
CausedByDrain.NORMAL);
201193
TimerData synchronizedProcessingTime2 =
202194
TimerData.of(
203195
NS1,
204196
new Instant(29),
205197
new Instant(29),
206198
TimeDomain.SYNCHRONIZED_PROCESSING_TIME,
207-
TimerData.CausedByDrain.NORMAL);
199+
CausedByDrain.NORMAL);
208200

209201
underTest.setTimer(processingTime1);
210202
underTest.setTimer(eventTime1);
@@ -239,18 +231,14 @@ public void testDeduplicate() throws Exception {
239231
InMemoryTimerInternals underTest = new InMemoryTimerInternals();
240232
TimerData eventTime =
241233
TimerData.of(
242-
NS1,
243-
new Instant(19),
244-
new Instant(19),
245-
TimeDomain.EVENT_TIME,
246-
TimerData.CausedByDrain.NORMAL);
234+
NS1, new Instant(19), new Instant(19), TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
247235
TimerData processingTime =
248236
TimerData.of(
249237
NS1,
250238
new Instant(19),
251239
new Instant(19),
252240
TimeDomain.PROCESSING_TIME,
253-
TimerData.CausedByDrain.NORMAL);
241+
CausedByDrain.NORMAL);
254242
underTest.setTimer(eventTime);
255243
underTest.setTimer(eventTime);
256244
underTest.setTimer(processingTime);

runners/core-java/src/test/java/org/apache/beam/runners/core/KeyedWorkItemCoderTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ public void testEncodeDecodeEqual() throws Exception {
4949
new Instant(500L),
5050
new Instant(500L),
5151
TimeDomain.EVENT_TIME,
52-
TimerData.CausedByDrain.NORMAL));
52+
CausedByDrain.NORMAL));
5353
Iterable<WindowedValue<Integer>> elements =
5454
ImmutableList.of(
5555
WindowedValues.valueInGlobalWindow(1),

runners/core-java/src/test/java/org/apache/beam/runners/core/ReduceFnTester.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -578,7 +578,7 @@ public void fireTimer(W window, Instant timestamp, TimeDomain domain) throws Exc
578578
timestamp,
579579
timestamp,
580580
domain,
581-
TimerData.CausedByDrain.NORMAL));
581+
CausedByDrain.NORMAL));
582582
runner.onTimers(timers);
583583
runner.persist();
584584
}
@@ -593,7 +593,7 @@ public void fireTimers(W window, TimestampedValue<TimeDomain>... timers) throws
593593
timer.getTimestamp(),
594594
timer.getTimestamp(),
595595
timer.getValue(),
596-
TimerData.CausedByDrain.NORMAL));
596+
CausedByDrain.NORMAL));
597597
}
598598
runner.onTimers(timerData);
599599
runner.persist();

runners/core-java/src/test/java/org/apache/beam/runners/core/SimpleDoFnRunnerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -702,7 +702,7 @@ public void onTimer(OnTimerContext context) {
702702
context.fireTimestamp(),
703703
context.timestamp(),
704704
context.timeDomain(),
705-
TimerData.CausedByDrain.NORMAL));
705+
CausedByDrain.NORMAL));
706706
}
707707
}
708708

runners/core-java/src/test/java/org/apache/beam/runners/core/SimplePushbackSideInputDoFnRunnerTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -318,7 +318,7 @@ public void testOnTimerCalled() {
318318
timestamp,
319319
timestamp,
320320
TimeDomain.EVENT_TIME,
321-
TimerData.CausedByDrain.NORMAL)));
321+
CausedByDrain.NORMAL)));
322322
}
323323

324324
private static class TestDoFnRunner<InputT, OutputT> implements DoFnRunner<InputT, OutputT> {
@@ -361,7 +361,7 @@ public <KeyT> void onTimer(
361361
timestamp,
362362
outputTimestamp,
363363
timeDomain,
364-
TimerData.CausedByDrain.NORMAL));
364+
CausedByDrain.NORMAL));
365365
}
366366

367367
@Override

runners/core-java/src/test/java/org/apache/beam/runners/core/TimerInternalsTest.java

Lines changed: 13 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ public void testTimerDataCoder() throws Exception {
4747
new Instant(0),
4848
new Instant(0),
4949
TimeDomain.EVENT_TIME,
50-
TimerData.CausedByDrain.NORMAL));
50+
CausedByDrain.NORMAL));
5151

5252
Coder<IntervalWindow> windowCoder = IntervalWindow.getCoder();
5353
CoderProperties.coderDecodeEncodeEqual(
@@ -59,7 +59,7 @@ windowCoder, new IntervalWindow(new Instant(0), new Instant(100))),
5959
new Instant(99),
6060
new Instant(99),
6161
TimeDomain.PROCESSING_TIME,
62-
TimerData.CausedByDrain.NORMAL));
62+
CausedByDrain.NORMAL));
6363
}
6464

6565
@Test
@@ -73,12 +73,7 @@ public void testCompareEqual() {
7373
StateNamespace namespace = StateNamespaces.global();
7474
TimerData timer =
7575
TimerData.of(
76-
"id",
77-
namespace,
78-
timestamp,
79-
timestamp,
80-
TimeDomain.EVENT_TIME,
81-
TimerData.CausedByDrain.NORMAL);
76+
"id", namespace, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
8277

8378
assertThat(
8479
timer,
@@ -89,7 +84,7 @@ public void testCompareEqual() {
8984
timestamp,
9085
timestamp,
9186
TimeDomain.EVENT_TIME,
92-
TimerData.CausedByDrain.NORMAL)));
87+
CausedByDrain.NORMAL)));
9388
}
9489

9590
@Test
@@ -100,18 +95,14 @@ public void testCompareByTimestamp() {
10095

10196
TimerData firstTimer =
10297
TimerData.of(
103-
namespace,
104-
firstTimestamp,
105-
firstTimestamp,
106-
TimeDomain.EVENT_TIME,
107-
TimerData.CausedByDrain.NORMAL);
98+
namespace, firstTimestamp, firstTimestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
10899
TimerData secondTimer =
109100
TimerData.of(
110101
namespace,
111102
secondTimestamp,
112103
secondTimestamp,
113104
TimeDomain.EVENT_TIME,
114-
TimerData.CausedByDrain.NORMAL);
105+
CausedByDrain.NORMAL);
115106

116107
assertThat(firstTimer, lessThan(secondTimer));
117108
}
@@ -122,22 +113,17 @@ public void testCompareByDomain() {
122113
StateNamespace namespace = StateNamespaces.global();
123114

124115
TimerData eventTimer =
125-
TimerData.of(
126-
namespace, timestamp, timestamp, TimeDomain.EVENT_TIME, TimerData.CausedByDrain.NORMAL);
116+
TimerData.of(namespace, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
127117
TimerData procTimer =
128118
TimerData.of(
129-
namespace,
130-
timestamp,
131-
timestamp,
132-
TimeDomain.PROCESSING_TIME,
133-
TimerData.CausedByDrain.NORMAL);
119+
namespace, timestamp, timestamp, TimeDomain.PROCESSING_TIME, CausedByDrain.NORMAL);
134120
TimerData synchronizedProcTimer =
135121
TimerData.of(
136122
namespace,
137123
timestamp,
138124
timestamp,
139125
TimeDomain.SYNCHRONIZED_PROCESSING_TIME,
140-
TimerData.CausedByDrain.NORMAL);
126+
CausedByDrain.NORMAL);
141127

142128
assertThat(eventTimer, lessThan(procTimer));
143129
assertThat(eventTimer, lessThan(synchronizedProcTimer));
@@ -156,18 +142,10 @@ public void testCompareByNamespace() {
156142

157143
TimerData secondEventTime =
158144
TimerData.of(
159-
firstWindowNs,
160-
timestamp,
161-
timestamp,
162-
TimeDomain.EVENT_TIME,
163-
TimerData.CausedByDrain.NORMAL);
145+
firstWindowNs, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
164146
TimerData thirdEventTime =
165147
TimerData.of(
166-
secondWindowNs,
167-
timestamp,
168-
timestamp,
169-
TimeDomain.EVENT_TIME,
170-
TimerData.CausedByDrain.NORMAL);
148+
secondWindowNs, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
171149

172150
assertThat(secondEventTime, lessThan(thirdEventTime));
173151
}
@@ -179,20 +157,10 @@ public void testCompareByTimerId() {
179157

180158
TimerData id0Timer =
181159
TimerData.of(
182-
"id0",
183-
namespace,
184-
timestamp,
185-
timestamp,
186-
TimeDomain.EVENT_TIME,
187-
TimerData.CausedByDrain.NORMAL);
160+
"id0", namespace, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
188161
TimerData id1Timer =
189162
TimerData.of(
190-
"id1",
191-
namespace,
192-
timestamp,
193-
timestamp,
194-
TimeDomain.EVENT_TIME,
195-
TimerData.CausedByDrain.NORMAL);
163+
"id1", namespace, timestamp, timestamp, TimeDomain.EVENT_TIME, CausedByDrain.NORMAL);
196164

197165
assertThat(id0Timer, lessThan(id1Timer));
198166
}

0 commit comments

Comments
 (0)