Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateInternals;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateReader;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateTagUtil;
import org.apache.beam.sdk.annotations.Internal;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.io.UnboundedSource;
Expand Down Expand Up @@ -772,6 +773,7 @@ public void start(
stateReader,
getWorkItem().getIsNewKey(),
cacheForKey.forFamily(stateFamily),
WindmillStateTagUtil.instance(),
scopedReadStateSupplier);

this.systemTimerInternals =
Expand All @@ -780,6 +782,7 @@ public void start(
WindmillNamespacePrefix.SYSTEM_NAMESPACE_PREFIX,
processingTime,
watermarks,
WindmillStateTagUtil.instance(),
td -> {});

this.userTimerInternals =
Expand All @@ -788,6 +791,7 @@ public void start(
WindmillNamespacePrefix.USER_NAMESPACE_PREFIX,
processingTime,
watermarks,
WindmillStateTagUtil.instance(),
this::onUserTimerModified);

this.cachedFiredSystemTimers = null;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,28 +17,30 @@
*/
package org.apache.beam.runners.dataflow.worker;

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

/**
* A prefix for a Windmill state or timer tag to separate user state and timers from system state
* and timers.
*/
enum WindmillNamespacePrefix {
@Internal
public enum WindmillNamespacePrefix {
USER_NAMESPACE_PREFIX {
@Override
ByteString byteString() {
public ByteString byteString() {
return USER_NAMESPACE_BYTESTRING;
}
},

SYSTEM_NAMESPACE_PREFIX {
@Override
ByteString byteString() {
public ByteString byteString() {
return SYSTEM_NAMESPACE_BYTESTRING;
}
};

abstract ByteString byteString();
public abstract ByteString byteString();

private static final ByteString USER_NAMESPACE_BYTESTRING = ByteString.copyFromUtf8("/u");
private static final ByteString SYSTEM_NAMESPACE_BYTESTRING = ByteString.copyFromUtf8("/s");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.beam.runners.dataflow.worker.streaming.Watermarks;
import org.apache.beam.runners.dataflow.worker.windmill.Windmill;
import org.apache.beam.runners.dataflow.worker.windmill.Windmill.Timer;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateTagUtil;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.state.TimeDomain;
import org.apache.beam.sdk.transforms.windowing.BoundedWindow;
Expand Down Expand Up @@ -59,7 +60,6 @@ class WindmillTimerInternals implements TimerInternals {
private static final Instant OUTPUT_TIMESTAMP_MAX_VALUE =
BoundedWindow.TIMESTAMP_MAX_VALUE.plus(Duration.millis(1));

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

public WindmillTimerInternals(
String stateFamily, // unique identifies a step
WindmillNamespacePrefix prefix, // partitions user and system namespaces into "/u" and "/s"
Instant processingTime,
Watermarks watermarks,
WindmillStateTagUtil windmillStateTagUtil,
Consumer<TimerData> onTimerModified) {
this.watermarks = watermarks;
this.processingTime = checkNotNull(processingTime);
this.stateFamily = stateFamily;
this.prefix = prefix;
this.windmillStateTagUtil = windmillStateTagUtil;
this.onTimerModified = onTimerModified;
}

public WindmillTimerInternals withPrefix(WindmillNamespacePrefix prefix) {
return new WindmillTimerInternals(
stateFamily, prefix, processingTime, watermarks, onTimerModified);
stateFamily, prefix, processingTime, watermarks, windmillStateTagUtil, onTimerModified);
}

@Override
Expand Down Expand Up @@ -211,7 +214,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
// Setting a timer, clear any prior hold and set to the new value
outputBuilder
.addWatermarkHoldsBuilder()
.setTag(timerHoldTag(prefix, timerData))
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
.setStateFamily(stateFamily)
.setReset(true)
.addTimestamps(
Expand All @@ -220,7 +223,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
// Clear the hold in case a previous iteration of this timer set one.
outputBuilder
.addWatermarkHoldsBuilder()
.setTag(timerHoldTag(prefix, timerData))
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
.setStateFamily(stateFamily)
.setReset(true);
}
Expand All @@ -235,7 +238,7 @@ public void persistTo(Windmill.WorkItemCommitRequest.Builder outputBuilder) {
// We are deleting timer; clear the hold
outputBuilder
.addWatermarkHoldsBuilder()
.setTag(timerHoldTag(prefix, timerData))
.setTag(windmillStateTagUtil.timerHoldTag(prefix, timerData))
.setStateFamily(stateFamily)
.setReset(true);
}
Expand Down Expand Up @@ -431,42 +434,6 @@ public static ByteString timerTag(WindmillNamespacePrefix prefix, TimerData time
return ByteString.copyFromUtf8(tagString);
}

/**
* Produce a state tag that is guaranteed to be unique for the given timer, to add a watermark
* hold that is only freed after the timer fires.
*/
public static ByteString timerHoldTag(WindmillNamespacePrefix prefix, TimerData timerData) {
String tagString;
if ("".equals(timerData.getTimerFamilyId())) {
tagString =
prefix.byteString().toStringUtf8()
+ // this never ends with a slash
TIMER_HOLD_PREFIX
+ // this never ends with a slash
timerData.getNamespace().stringKey()
+ // this must begin and end with a slash
'+'
+ timerData.getTimerId() // this is arbitrary; currently unescaped
;
} else {
tagString =
prefix.byteString().toStringUtf8()
+ // this never ends with a slash
TIMER_HOLD_PREFIX
+ // this never ends with a slash
timerData.getNamespace().stringKey()
+ // this must begin and end with a slash
'+'
+ timerData.getTimerId()
+ // this is arbitrary; currently unescaped
'+'
+ timerData.getTimerFamilyId() // use to differentiate same timerId in different
// timerMap
;
}
return ByteString.copyFromUtf8(tagString);
}

@VisibleForTesting
static Timer.Type timerType(TimeDomain domain) {
switch (domain) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.beam.runners.core.StateTable;
import org.apache.beam.runners.core.StateTag;
import org.apache.beam.runners.core.StateTags;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache.ForKeyAndFamily;
import org.apache.beam.sdk.coders.BooleanCoder;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.state.*;
Expand All @@ -43,6 +44,7 @@ final class CachingStateTable extends StateTable {
private final @Nullable StateTable derivedStateTable;
private final boolean isNewKey;
private final boolean mapStateViaMultimapState;
private final WindmillStateTagUtil windmillStateTagUtil;

private CachingStateTable(Builder builder) {
this.stateFamily = builder.stateFamily;
Expand All @@ -53,7 +55,7 @@ private CachingStateTable(Builder builder) {
this.scopedReadStateSupplier = builder.scopedReadStateSupplier;
this.derivedStateTable = builder.derivedStateTable;
this.mapStateViaMultimapState = builder.mapStateViaMultimapState;

this.windmillStateTagUtil = builder.windmillStateTagUtil;
if (this.isSystemTable) {
Preconditions.checkState(derivedStateTable == null);
} else {
Expand All @@ -64,11 +66,12 @@ private CachingStateTable(Builder builder) {
static CachingStateTable.Builder builder(
String stateFamily,
WindmillStateReader reader,
WindmillStateCache.ForKeyAndFamily cache,
ForKeyAndFamily cache,
boolean isNewKey,
Supplier<Closeable> scopedReadStateSupplier) {
Supplier<Closeable> scopedReadStateSupplier,
WindmillStateTagUtil windmillStateTagUtil) {
return new CachingStateTable.Builder(
stateFamily, reader, cache, scopedReadStateSupplier, isNewKey);
stateFamily, reader, cache, scopedReadStateSupplier, isNewKey, windmillStateTagUtil);
}

@Override
Expand All @@ -89,7 +92,12 @@ public <T> BagState<T> bindBag(StateTag<BagState<T>> address, Coder<T> elemCoder
.orElseGet(
() ->
new WindmillBag<>(
namespace, resolvedAddress, stateFamily, elemCoder, isNewKey));
namespace,
resolvedAddress,
stateFamily,
elemCoder,
isNewKey,
windmillStateTagUtil));

result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand Down Expand Up @@ -122,7 +130,13 @@ public <KeyT, ValueT> AbstractWindmillMap<KeyT, ValueT> bindMap(
.orElseGet(
() ->
new WindmillMap<>(
namespace, spec, stateFamily, keyCoder, valueCoder, isNewKey));
namespace,
spec,
stateFamily,
keyCoder,
valueCoder,
isNewKey,
windmillStateTagUtil));
}
result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand All @@ -140,7 +154,13 @@ public <KeyT, ValueT> WindmillMultimap<KeyT, ValueT> bindMultimap(
.orElseGet(
() ->
new WindmillMultimap<>(
namespace, spec, stateFamily, keyCoder, valueCoder, isNewKey));
namespace,
spec,
stateFamily,
keyCoder,
valueCoder,
isNewKey,
windmillStateTagUtil));
result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
}
Expand All @@ -162,7 +182,8 @@ public <T> OrderedListState<T> bindOrderedList(
specOrInternalTag,
stateFamily,
elemCoder,
isNewKey));
isNewKey,
windmillStateTagUtil));

result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand All @@ -180,7 +201,12 @@ public WatermarkHoldState bindWatermark(
.orElseGet(
() ->
new WindmillWatermarkHold(
namespace, address, stateFamily, timestampCombiner, isNewKey));
namespace,
address,
stateFamily,
timestampCombiner,
isNewKey,
windmillStateTagUtil));

result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand All @@ -202,7 +228,8 @@ public <InputT, AccumT, OutputT> CombiningState<InputT, AccumT, OutputT> bindCom
accumCoder,
combineFn,
cache,
isNewKey);
isNewKey,
windmillStateTagUtil);

result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand All @@ -229,7 +256,12 @@ public <T> ValueState<T> bindValue(StateTag<ValueState<T>> address, Coder<T> cod
.orElseGet(
() ->
new WindmillValue<>(
namespace, addressOrInternalTag, stateFamily, coder, isNewKey));
namespace,
addressOrInternalTag,
stateFamily,
coder,
isNewKey,
windmillStateTagUtil));

result.initializeForWorkItem(reader, scopedReadStateSupplier);
return result;
Expand All @@ -247,23 +279,26 @@ static class Builder {
private final WindmillStateCache.ForKeyAndFamily cache;
private final Supplier<Closeable> scopedReadStateSupplier;
private final boolean isNewKey;
private final WindmillStateTagUtil windmillStateTagUtil;
private boolean isSystemTable;
private @Nullable StateTable derivedStateTable;
private boolean mapStateViaMultimapState = false;

private Builder(
String stateFamily,
WindmillStateReader reader,
WindmillStateCache.ForKeyAndFamily cache,
ForKeyAndFamily cache,
Supplier<Closeable> scopedReadStateSupplier,
boolean isNewKey) {
boolean isNewKey,
WindmillStateTagUtil windmillStateTagUtil) {
this.stateFamily = stateFamily;
this.reader = reader;
this.cache = cache;
this.scopedReadStateSupplier = scopedReadStateSupplier;
this.isNewKey = isNewKey;
this.isSystemTable = true;
this.derivedStateTable = null;
this.windmillStateTagUtil = windmillStateTagUtil;
}

Builder withDerivedState(StateTable derivedStateTable) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,10 +63,11 @@ public class WindmillBag<T> extends SimpleWindmillState implements BagState<T> {
StateTag<BagState<T>> address,
String stateFamily,
Coder<T> elemCoder,
boolean isNewKey) {
boolean isNewKey,
WindmillStateTagUtil windmillStateTagUtil) {
this.namespace = namespace;
this.address = address;
this.stateKey = WindmillStateUtil.encodeKey(namespace, address);
this.stateKey = windmillStateTagUtil.encodeKey(namespace, address);
this.stateFamily = stateFamily;
this.elemCoder = elemCoder;
if (isNewKey) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,13 @@
import org.apache.beam.runners.core.StateTag;
import org.apache.beam.runners.core.StateTags;
import org.apache.beam.runners.dataflow.worker.windmill.Windmill;
import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache.ForKeyAndFamily;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.state.BagState;
import org.apache.beam.sdk.state.CombiningState;
import org.apache.beam.sdk.state.ReadableState;
import org.apache.beam.sdk.transforms.Combine;
import org.apache.beam.sdk.transforms.Combine.CombineFn;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Supplier;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.Iterables;

Expand All @@ -54,9 +56,10 @@ class WindmillCombiningState<InputT, AccumT, OutputT> extends WindmillState
StateTag<CombiningState<InputT, AccumT, OutputT>> address,
String stateFamily,
Coder<AccumT> accumCoder,
Combine.CombineFn<InputT, AccumT, OutputT> combineFn,
WindmillStateCache.ForKeyAndFamily cache,
boolean isNewKey) {
CombineFn<InputT, AccumT, OutputT> combineFn,
ForKeyAndFamily cache,
boolean isNewKey,
WindmillStateTagUtil windmillStateTagUtil) {
StateTag<BagState<AccumT>> internalBagAddress = StateTags.convertToBagTagInternal(address);
this.bag =
cache
Expand All @@ -65,7 +68,12 @@ class WindmillCombiningState<InputT, AccumT, OutputT> extends WindmillState
.orElseGet(
() ->
new WindmillBag<>(
namespace, internalBagAddress, stateFamily, accumCoder, isNewKey));
namespace,
internalBagAddress,
stateFamily,
accumCoder,
isNewKey,
windmillStateTagUtil));

this.combineFn = combineFn;
this.localAdditionsAccumulator = combineFn.createAccumulator();
Expand Down
Loading
Loading