Skip to content

Commit 65ee225

Browse files
authored
[Dataflow Streaming] Move functionality creating windmill tags to a common class. (#36283)
This is a prep to introduce new windmill tag encodings. No functionality change.
1 parent c5a6189 commit 65ee225

15 files changed

Lines changed: 162 additions & 92 deletions

File tree

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@
6161
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache;
6262
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateInternals;
6363
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateReader;
64+
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateTagUtil;
6465
import org.apache.beam.sdk.annotations.Internal;
6566
import org.apache.beam.sdk.coders.Coder;
6667
import org.apache.beam.sdk.io.UnboundedSource;
@@ -772,6 +773,7 @@ public void start(
772773
stateReader,
773774
getWorkItem().getIsNewKey(),
774775
cacheForKey.forFamily(stateFamily),
776+
WindmillStateTagUtil.instance(),
775777
scopedReadStateSupplier);
776778

777779
this.systemTimerInternals =
@@ -780,6 +782,7 @@ public void start(
780782
WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX,
781783
processingTime,
782784
watermarks,
785+
WindmillStateTagUtil.instance(),
783786
td -> {});
784787

785788
this.userTimerInternals =
@@ -788,6 +791,7 @@ public void start(
788791
WindmillNamespacePrefix.USER_NAMESPACE_PREFIX,
789792
processingTime,
790793
watermarks,
794+
WindmillStateTagUtil.instance(),
791795
this::onUserTimerModified);
792796

793797
this.cachedFiredSystemTimers = null;

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillNamespacePrefix.java

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,28 +17,30 @@
1717
*/
1818
package org.apache.beam.runners.dataflow.worker;
1919

20+
import org.apache.beam.sdk.annotations.Internal;
2021
import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString;
2122

2223
/**
2324
* A prefix for a Windmill state or timer tag to separate user state and timers from system state
2425
* and timers.
2526
*/
26-
enum WindmillNamespacePrefix {
27+
@Internal
28+
public enum WindmillNamespacePrefix {
2729
USER_NAMESPACE_PREFIX {
2830
@Override
29-
ByteString byteString() {
31+
public ByteString byteString() {
3032
return USER_NAMESPACE_BYTESTRING;
3133
}
3234
},
3335

3436
SYSTEM_NAMESPACE_PREFIX {
3537
@Override
36-
ByteString byteString() {
38+
public ByteString byteString() {
3739
return SYSTEM_NAMESPACE_BYTESTRING;
3840
}
3941
};
4042

41-
abstract ByteString byteString();
43+
public abstract ByteString byteString();
4244

4345
private static final ByteString USER_NAMESPACE_BYTESTRING = ByteString.copyFromUtf8("/u");
4446
private static final ByteString SYSTEM_NAMESPACE_BYTESTRING = ByteString.copyFromUtf8("/s");

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java

Lines changed: 8 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import org.apache.beam.runners.dataflow.worker.streaming.Watermarks;
3333
import org.apache.beam.runners.dataflow.worker.windmill.Windmill;
3434
import org.apache.beam.runners.dataflow.worker.windmill.Windmill.Timer;
35+
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateTagUtil;
3536
import org.apache.beam.sdk.coders.Coder;
3637
import org.apache.beam.sdk.state.TimeDomain;
3738
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
@@ -59,7 +60,6 @@ class WindmillTimerInternals implements TimerInternals {
5960
private static final Instant OUTPUT_TIMESTAMP_MAX_VALUE =
6061
BoundedWindow.TIMESTAMP_MAX_VALUE.plus(Duration.millis(1));
6162

62-
private static final String TIMER_HOLD_PREFIX = "/h";
6363
// Map from timer id to its TimerData. If it is to be deleted, we still need
6464
// its time domain here. Note that TimerData is unique per ID and namespace,
6565
// though technically in Windmill this is only enforced per ID and namespace
@@ -74,23 +74,26 @@ class WindmillTimerInternals implements TimerInternals {
7474
private final String stateFamily;
7575
private final WindmillNamespacePrefix prefix;
7676
private final Consumer<TimerData> onTimerModified;
77+
private final WindmillStateTagUtil windmillStateTagUtil;
7778

7879
public WindmillTimerInternals(
7980
String stateFamily, // unique identifies a step
8081
WindmillNamespacePrefix prefix, // partitions user and system namespaces into "/u" and "/s"
8182
Instant processingTime,
8283
Watermarks watermarks,
84+
WindmillStateTagUtil windmillStateTagUtil,
8385
Consumer<TimerData> onTimerModified) {
8486
this.watermarks = watermarks;
8587
this.processingTime = checkNotNull(processingTime);
8688
this.stateFamily = stateFamily;
8789
this.prefix = prefix;
90+
this.windmillStateTagUtil = windmillStateTagUtil;
8891
this.onTimerModified = onTimerModified;
8992
}
9093

9194
public WindmillTimerInternals withPrefix(WindmillNamespacePrefix prefix) {
9295
return new WindmillTimerInternals(
93-
stateFamily, prefix, processingTime, watermarks, onTimerModified);
96+
stateFamily, prefix, processingTime, watermarks, windmillStateTagUtil, onTimerModified);
9497
}
9598

9699
@Override
@@ -211,7 +214,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
211214
// Setting a timer, clear any prior hold and set to the new value
212215
outputBuilder
213216
.addWatermarkHoldsBuilder()
214-
.setTag(timerHoldTag(prefix, timerData))
217+
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
215218
.setStateFamily(stateFamily)
216219
.setReset(true)
217220
.addTimestamps(
@@ -220,7 +223,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
220223
// Clear the hold in case a previous iteration of this timer set one.
221224
outputBuilder
222225
.addWatermarkHoldsBuilder()
223-
.setTag(timerHoldTag(prefix, timerData))
226+
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
224227
.setStateFamily(stateFamily)
225228
.setReset(true);
226229
}
@@ -235,7 +238,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
235238
// We are deleting timer; clear the hold
236239
outputBuilder
237240
.addWatermarkHoldsBuilder()
238-
.setTag(timerHoldTag(prefix, timerData))
241+
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
239242
.setStateFamily(stateFamily)
240243
.setReset(true);
241244
}
@@ -431,42 +434,6 @@ public static ByteString timerTag(WindmillNamespacePrefix prefix, TimerData time
431434
return ByteString.copyFromUtf8(tagString);
432435
}
433436

434-
/**
435-
* Produce a state tag that is guaranteed to be unique for the given timer, to add a watermark
436-
* hold that is only freed after the timer fires.
437-
*/
438-
public static ByteString timerHoldTag(WindmillNamespacePrefix prefix, TimerData timerData) {
439-
String tagString;
440-
if ("".equals(timerData.getTimerFamilyId())) {
441-
tagString =
442-
prefix.byteString().toStringUtf8()
443-
+ // this never ends with a slash
444-
TIMER_HOLD_PREFIX
445-
+ // this never ends with a slash
446-
timerData.getNamespace().stringKey()
447-
+ // this must begin and end with a slash
448-
'+'
449-
+ timerData.getTimerId() // this is arbitrary; currently unescaped
450-
;
451-
} else {
452-
tagString =
453-
prefix.byteString().toStringUtf8()
454-
+ // this never ends with a slash
455-
TIMER_HOLD_PREFIX
456-
+ // this never ends with a slash
457-
timerData.getNamespace().stringKey()
458-
+ // this must begin and end with a slash
459-
'+'
460-
+ timerData.getTimerId()
461-
+ // this is arbitrary; currently unescaped
462-
'+'
463-
+ timerData.getTimerFamilyId() // use to differentiate same timerId in different
464-
// timerMap
465-
;
466-
}
467-
return ByteString.copyFromUtf8(tagString);
468-
}
469-
470437
@VisibleForTesting
471438
static Timer.Type timerType(TimeDomain domain) {
472439
switch (domain) {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/CachingStateTable.java

Lines changed: 48 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import org.apache.beam.runners.core.StateTable;
2525
import org.apache.beam.runners.core.StateTag;
2626
import org.apache.beam.runners.core.StateTags;
27+
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache.ForKeyAndFamily;
2728
import org.apache.beam.sdk.coders.BooleanCoder;
2829
import org.apache.beam.sdk.coders.Coder;
2930
import org.apache.beam.sdk.state.*;
@@ -43,6 +44,7 @@ final class CachingStateTable extends StateTable {
4344
private final @Nullable StateTable derivedStateTable;
4445
private final boolean isNewKey;
4546
private final boolean mapStateViaMultimapState;
47+
private final WindmillStateTagUtil windmillStateTagUtil;
4648

4749
private CachingStateTable(Builder builder) {
4850
this.stateFamily = builder.stateFamily;
@@ -53,7 +55,7 @@ private CachingStateTable(Builder builder) {
5355
this.scopedReadStateSupplier = builder.scopedReadStateSupplier;
5456
this.derivedStateTable = builder.derivedStateTable;
5557
this.mapStateViaMultimapState = builder.mapStateViaMultimapState;
56-
58+
this.windmillStateTagUtil = builder.windmillStateTagUtil;
5759
if (this.isSystemTable) {
5860
Preconditions.checkState(derivedStateTable == null);
5961
} else {
@@ -64,11 +66,12 @@ private CachingStateTable(Builder builder) {
6466
static CachingStateTable.Builder builder(
6567
String stateFamily,
6668
WindmillStateReader reader,
67-
WindmillStateCache.ForKeyAndFamily cache,
69+
ForKeyAndFamily cache,
6870
boolean isNewKey,
69-
Supplier<Closeable> scopedReadStateSupplier) {
71+
Supplier<Closeable> scopedReadStateSupplier,
72+
WindmillStateTagUtil windmillStateTagUtil) {
7073
return new CachingStateTable.Builder(
71-
stateFamily, reader, cache, scopedReadStateSupplier, isNewKey);
74+
stateFamily, reader, cache, scopedReadStateSupplier, isNewKey, windmillStateTagUtil);
7275
}
7376

7477
@Override
@@ -89,7 +92,12 @@ public <T> BagState<T> bindBag(StateTag<BagState<T>> address, Coder<T> elemCoder
8992
.orElseGet(
9093
() ->
9194
new WindmillBag<>(
92-
namespace, resolvedAddress, stateFamily, elemCoder, isNewKey));
95+
namespace,
96+
resolvedAddress,
97+
stateFamily,
98+
elemCoder,
99+
isNewKey,
100+
windmillStateTagUtil));
93101

94102
result.initializeForWorkItem(reader, scopedReadStateSupplier);
95103
return result;
@@ -122,7 +130,13 @@ public <KeyT, ValueT> AbstractWindmillMap<KeyT, ValueT> bindMap(
122130
.orElseGet(
123131
() ->
124132
new WindmillMap<>(
125-
namespace, spec, stateFamily, keyCoder, valueCoder, isNewKey));
133+
namespace,
134+
spec,
135+
stateFamily,
136+
keyCoder,
137+
valueCoder,
138+
isNewKey,
139+
windmillStateTagUtil));
126140
}
127141
result.initializeForWorkItem(reader, scopedReadStateSupplier);
128142
return result;
@@ -140,7 +154,13 @@ public <KeyT, ValueT> WindmillMultimap<KeyT, ValueT> bindMultimap(
140154
.orElseGet(
141155
() ->
142156
new WindmillMultimap<>(
143-
namespace, spec, stateFamily, keyCoder, valueCoder, isNewKey));
157+
namespace,
158+
spec,
159+
stateFamily,
160+
keyCoder,
161+
valueCoder,
162+
isNewKey,
163+
windmillStateTagUtil));
144164
result.initializeForWorkItem(reader, scopedReadStateSupplier);
145165
return result;
146166
}
@@ -162,7 +182,8 @@ public <T> OrderedListState<T> bindOrderedList(
162182
specOrInternalTag,
163183
stateFamily,
164184
elemCoder,
165-
isNewKey));
185+
isNewKey,
186+
windmillStateTagUtil));
166187

167188
result.initializeForWorkItem(reader, scopedReadStateSupplier);
168189
return result;
@@ -180,7 +201,12 @@ public WatermarkHoldState bindWatermark(
180201
.orElseGet(
181202
() ->
182203
new WindmillWatermarkHold(
183-
namespace, address, stateFamily, timestampCombiner, isNewKey));
204+
namespace,
205+
address,
206+
stateFamily,
207+
timestampCombiner,
208+
isNewKey,
209+
windmillStateTagUtil));
184210

185211
result.initializeForWorkItem(reader, scopedReadStateSupplier);
186212
return result;
@@ -202,7 +228,8 @@ public <InputT, AccumT, OutputT> CombiningState<InputT, AccumT, OutputT> bindCom
202228
accumCoder,
203229
combineFn,
204230
cache,
205-
isNewKey);
231+
isNewKey,
232+
windmillStateTagUtil);
206233

207234
result.initializeForWorkItem(reader, scopedReadStateSupplier);
208235
return result;
@@ -229,7 +256,12 @@ public <T> ValueState<T> bindValue(StateTag<ValueState<T>> address, Coder<T> cod
229256
.orElseGet(
230257
() ->
231258
new WindmillValue<>(
232-
namespace, addressOrInternalTag, stateFamily, coder, isNewKey));
259+
namespace,
260+
addressOrInternalTag,
261+
stateFamily,
262+
coder,
263+
isNewKey,
264+
windmillStateTagUtil));
233265

234266
result.initializeForWorkItem(reader, scopedReadStateSupplier);
235267
return result;
@@ -247,23 +279,26 @@ static class Builder {
247279
private final WindmillStateCache.ForKeyAndFamily cache;
248280
private final Supplier<Closeable> scopedReadStateSupplier;
249281
private final boolean isNewKey;
282+
private final WindmillStateTagUtil windmillStateTagUtil;
250283
private boolean isSystemTable;
251284
private @Nullable StateTable derivedStateTable;
252285
private boolean mapStateViaMultimapState = false;
253286

254287
private Builder(
255288
String stateFamily,
256289
WindmillStateReader reader,
257-
WindmillStateCache.ForKeyAndFamily cache,
290+
ForKeyAndFamily cache,
258291
Supplier<Closeable> scopedReadStateSupplier,
259-
boolean isNewKey) {
292+
boolean isNewKey,
293+
WindmillStateTagUtil windmillStateTagUtil) {
260294
this.stateFamily = stateFamily;
261295
this.reader = reader;
262296
this.cache = cache;
263297
this.scopedReadStateSupplier = scopedReadStateSupplier;
264298
this.isNewKey = isNewKey;
265299
this.isSystemTable = true;
266300
this.derivedStateTable = null;
301+
this.windmillStateTagUtil = windmillStateTagUtil;
267302
}
268303

269304
Builder withDerivedState(StateTable derivedStateTable) {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillBag.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -63,10 +63,11 @@ public class WindmillBag<T> extends SimpleWindmillState implements BagState<T> {
6363
StateTag<BagState<T>> address,
6464
String stateFamily,
6565
Coder<T> elemCoder,
66-
boolean isNewKey) {
66+
boolean isNewKey,
67+
WindmillStateTagUtil windmillStateTagUtil) {
6768
this.namespace = namespace;
6869
this.address = address;
69-
this.stateKey = WindmillStateUtil.encodeKey(namespace, address);
70+
this.stateKey = windmillStateTagUtil.encodeKey(namespace, address);
7071
this.stateFamily = stateFamily;
7172
this.elemCoder = elemCoder;
7273
if (isNewKey) {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillCombiningState.java

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,13 @@
2727
import org.apache.beam.runners.core.StateTag;
2828
import org.apache.beam.runners.core.StateTags;
2929
import org.apache.beam.runners.dataflow.worker.windmill.Windmill;
30+
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache.ForKeyAndFamily;
3031
import org.apache.beam.sdk.coders.Coder;
3132
import org.apache.beam.sdk.state.BagState;
3233
import org.apache.beam.sdk.state.CombiningState;
3334
import org.apache.beam.sdk.state.ReadableState;
3435
import org.apache.beam.sdk.transforms.Combine;
36+
import org.apache.beam.sdk.transforms.Combine.CombineFn;
3537
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Supplier;
3638
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;
3739

@@ -54,9 +56,10 @@ class WindmillCombiningState<InputT, AccumT, OutputT> extends WindmillState
5456
StateTag<CombiningState<InputT, AccumT, OutputT>> address,
5557
String stateFamily,
5658
Coder<AccumT> accumCoder,
57-
Combine.CombineFn<InputT, AccumT, OutputT> combineFn,
58-
WindmillStateCache.ForKeyAndFamily cache,
59-
boolean isNewKey) {
59+
CombineFn<InputT, AccumT, OutputT> combineFn,
60+
ForKeyAndFamily cache,
61+
boolean isNewKey,
62+
WindmillStateTagUtil windmillStateTagUtil) {
6063
StateTag<BagState<AccumT>> internalBagAddress = StateTags.convertToBagTagInternal(address);
6164
this.bag =
6265
cache
@@ -65,7 +68,12 @@ class WindmillCombiningState<InputT, AccumT, OutputT> extends WindmillState
6568
.orElseGet(
6669
() ->
6770
new WindmillBag<>(
68-
namespace, internalBagAddress, stateFamily, accumCoder, isNewKey));
71+
namespace,
72+
internalBagAddress,
73+
stateFamily,
74+
accumCoder,
75+
isNewKey,
76+
windmillStateTagUtil));
6977

7078
this.combineFn = combineFn;
7179
this.localAdditionsAccumulator = combineFn.createAccumulator();

0 commit comments

Comments
 (0)