Skip to content

Commit b9a8972

Browse files
authored
External metadata for streaming runner v1 changes (#36373)
* encode empty metadata within WindowedValue * add todos for future otel metadata * add todos for future otel metadata * spotless * fix tests * fix * fix * fix bugs * fix bugs
1 parent d54a661 commit b9a8972

16 files changed

Lines changed: 212 additions & 27 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -135,6 +135,7 @@ public Duration getAllowedTimestampSkew() {
135135
public void processElement(
136136
@Element KV<K, ValueInSingleWindow<V>> kv,
137137
OutputReceiver<KV<K, V>> outputReceiver) {
138+
// todo #33176 specify additional metadata in the future
138139
outputReceiver
139140
.builder(KV.of(kv.getKey(), kv.getValue().getValue()))
140141
.setTimestamp(kv.getValue().getTimestamp())

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,7 @@
112112
import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics;
113113
import org.apache.beam.sdk.metrics.MetricsEnvironment;
114114
import org.apache.beam.sdk.util.construction.CoderTranslation;
115+
import org.apache.beam.sdk.values.WindowedValues;
115116
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
116117
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
117118
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.cache.CacheStats;
@@ -165,6 +166,7 @@ public final class StreamingDataflowWorker {
165166
private static final Random CLIENT_ID_GENERATOR = new Random();
166167
private static final String CHANNELZ_PATH = "/channelz";
167168
private static final String BEAM_FN_API_EXPERIMENT = "beam_fn_api";
169+
private static final String ELEMENT_METADATA_SUPPORTED_EXPERIMENT = "element_metadata_supported";
168170
private static final String STREAMING_ENGINE_USE_JOB_SETTINGS_FOR_HEARTBEAT_POOL_EXPERIMENT =
169171
"streaming_engine_use_job_settings_for_heartbeat_pool";
170172
// Experiment make the monitor within BoundedQueueExecutor fair
@@ -985,6 +987,9 @@ public static void main(String[] args) throws Exception {
985987
validateWorkerOptions(options);
986988

987989
CoderTranslation.verifyModelCodersRegistered();
990+
if (DataflowRunner.hasExperiment(options, ELEMENT_METADATA_SUPPORTED_EXPERIMENT)) {
991+
WindowedValues.FullWindowedValueCoder.setMetadataSupported();
992+
}
988993

989994
LOG.debug("Creating StreamingDataflowWorker from options: {}", options);
990995
StreamingDataflowWorker worker = StreamingDataflowWorker.fromOptions(options);

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,9 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
117117
Collection<? extends BoundedWindow> windows =
118118
WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata());
119119
PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata());
120+
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
121+
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
122+
}
120123
if (valueCoder instanceof KvCoder) {
121124
KvCoder<?, ?> kvCoder = (KvCoder<?, ?>) valueCoder;
122125
InputStream key = context.getSerializedKey().newInput();
@@ -125,9 +128,11 @@ protected WindowedValue<T> decodeMessage(Windmill.Message message) throws IOExce
125128
@SuppressWarnings("unchecked")
126129
T result =
127130
(T) KV.of(decode(kvCoder.getKeyCoder(), key), decode(kvCoder.getValueCoder(), data));
131+
// todo #33176 propagate metadata to windowed value
128132
return WindowedValues.of(result, timestampMillis, windows, paneInfo);
129133
} else {
130134
notifyElementRead(data.available() + metadata.available());
135+
// todo #33176 propagate metadata to windowed value
131136
return WindowedValues.of(decode(valueCoder, data), timestampMillis, windows, paneInfo);
132137
}
133138
}

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

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -108,9 +108,12 @@ public Iterable<WindowedValue<ElemT>> elementsIterable() {
108108
Collection<? extends BoundedWindow> windows =
109109
WindmillSink.decodeMetadataWindows(windowsCoder, message.getMetadata());
110110
PaneInfo paneInfo = WindmillSink.decodeMetadataPane(message.getMetadata());
111-
111+
if (WindowedValues.WindowedValueCoder.isMetadataSupported()) {
112+
WindmillSink.decodeAdditionalMetadata(windowsCoder, message.getMetadata());
113+
}
112114
InputStream inputStream = message.getData().newInput();
113115
ElemT value = valueCoder.decode(inputStream, Coder.Context.OUTER);
116+
// todo #33176 specify additional metadata in the future
114117
return WindowedValues.of(value, timestamp, windows, paneInfo);
115118
} catch (IOException e) {
116119
throw new RuntimeException(e);

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

Lines changed: 39 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,11 @@
2626
import java.util.Collection;
2727
import java.util.HashMap;
2828
import java.util.Map;
29+
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
2930
import org.apache.beam.runners.dataflow.util.CloudObject;
3031
import org.apache.beam.runners.dataflow.worker.util.common.worker.Sink;
3132
import org.apache.beam.runners.dataflow.worker.windmill.Windmill;
33+
import org.apache.beam.sdk.coders.ByteArrayCoder;
3234
import org.apache.beam.sdk.coders.Coder;
3335
import org.apache.beam.sdk.coders.KvCoder;
3436
import org.apache.beam.sdk.options.PipelineOptions;
@@ -40,6 +42,7 @@
4042
import org.apache.beam.sdk.values.ValueWithRecordId;
4143
import org.apache.beam.sdk.values.ValueWithRecordId.ValueWithRecordIdCoder;
4244
import org.apache.beam.sdk.values.WindowedValue;
45+
import org.apache.beam.sdk.values.WindowedValues;
4346
import org.apache.beam.sdk.values.WindowedValues.FullWindowedValueCoder;
4447
import org.apache.beam.vendor.grpc.v1p69p0.com.google.protobuf.ByteString;
4548
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
@@ -75,11 +78,20 @@ private static ByteString encodeMetadata(
7578
ByteStringOutputStream stream,
7679
Coder<Collection<? extends BoundedWindow>> windowsCoder,
7780
Collection<? extends BoundedWindow> windows,
78-
PaneInfo paneInfo)
81+
PaneInfo paneInfo,
82+
BeamFnApi.Elements.ElementMetadata metadata)
7983
throws IOException {
8084
try {
81-
PaneInfoCoder.INSTANCE.encode(paneInfo, stream);
82-
windowsCoder.encode(windows, stream, Coder.Context.OUTER);
85+
// element metadata is behind the experiment
86+
boolean elementMetadata = WindowedValues.WindowedValueCoder.isMetadataSupported();
87+
if (elementMetadata) {
88+
PaneInfoCoder.INSTANCE.encode(paneInfo.withElementMetadata(true), stream);
89+
windowsCoder.encode(windows, stream);
90+
ByteArrayCoder.of().encode(metadata.toByteArray(), stream, Coder.Context.OUTER);
91+
} else {
92+
PaneInfoCoder.INSTANCE.encode(paneInfo, stream);
93+
windowsCoder.encode(windows, stream, Coder.Context.OUTER);
94+
}
8395
return stream.toByteStringAndReset();
8496
} catch (Exception e) {
8597
stream.reset();
@@ -90,23 +102,39 @@ private static ByteString encodeMetadata(
90102
public static ByteString encodeMetadata(
91103
Coder<Collection<? extends BoundedWindow>> windowsCoder,
92104
Collection<? extends BoundedWindow> windows,
93-
PaneInfo paneInfo)
105+
PaneInfo paneInfo,
106+
BeamFnApi.Elements.ElementMetadata metadata)
94107
throws IOException {
95108
ByteStringOutputStream stream = new ByteStringOutputStream();
96-
return encodeMetadata(stream, windowsCoder, windows, paneInfo);
109+
return encodeMetadata(stream, windowsCoder, windows, paneInfo, metadata);
97110
}
98111

99112
public static PaneInfo decodeMetadataPane(ByteString metadata) throws IOException {
100113
InputStream inStream = metadata.newInput();
101114
return PaneInfoCoder.INSTANCE.decode(inStream);
102115
}
103116

117+
public static BeamFnApi.Elements.ElementMetadata decodeAdditionalMetadata(
118+
Coder<Collection<? extends BoundedWindow>> windowsCoder, ByteString metadata)
119+
throws IOException {
120+
InputStream inStream = metadata.newInput();
121+
PaneInfo paneInfo = PaneInfoCoder.INSTANCE.decode(inStream);
122+
windowsCoder.decode(inStream);
123+
if (paneInfo.isElementMetadata()) {
124+
return BeamFnApi.Elements.ElementMetadata.parseFrom(
125+
ByteArrayCoder.of().decode(inStream, Coder.Context.OUTER));
126+
} else {
127+
// empty
128+
return BeamFnApi.Elements.ElementMetadata.newBuilder().build();
129+
}
130+
}
131+
104132
public static Collection<? extends BoundedWindow> decodeMetadataWindows(
105133
Coder<Collection<? extends BoundedWindow>> windowsCoder, ByteString metadata)
106134
throws IOException {
107135
InputStream inStream = metadata.newInput();
108136
PaneInfoCoder.INSTANCE.decode(inStream);
109-
return windowsCoder.decode(inStream, Coder.Context.OUTER);
137+
return windowsCoder.decode(inStream);
110138
}
111139

112140
/** A {@link SinkFactory.Registrar} for windmill sinks. */
@@ -184,8 +212,12 @@ private <EncodeT> ByteString encode(Coder<EncodeT> coder, EncodeT object) throws
184212
public long add(WindowedValue<T> data) throws IOException {
185213
ByteString key, value;
186214
ByteString id = ByteString.EMPTY;
215+
// todo #33176 specify additional metadata in the future
216+
BeamFnApi.Elements.ElementMetadata additionalMetadata =
217+
BeamFnApi.Elements.ElementMetadata.newBuilder().build();
187218
ByteString metadata =
188-
encodeMetadata(stream, windowsCoder, data.getWindows(), data.getPaneInfo());
219+
encodeMetadata(
220+
stream, windowsCoder, data.getWindows(), data.getPaneInfo(), additionalMetadata);
189221
if (valueCoder instanceof KvCoder) {
190222
KvCoder kvCoder = (KvCoder) valueCoder;
191223
KV kv = (KV) data.getValue();

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowFnsTest.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
import java.util.Arrays;
3030
import java.util.Collection;
3131
import java.util.List;
32+
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
3233
import org.apache.beam.runners.core.DoFnRunner;
3334
import org.apache.beam.runners.core.DoFnRunners;
3435
import org.apache.beam.runners.core.InMemoryStateInternals;
@@ -176,7 +177,12 @@ private <V> void addElement(
176177
valueCoder.encode(value, dataOutput, Context.OUTER);
177178
messageBundle
178179
.addMessagesBuilder()
179-
.setMetadata(WindmillSink.encodeMetadata(windowsCoder, windows, PaneInfo.NO_FIRING))
180+
.setMetadata(
181+
WindmillSink.encodeMetadata(
182+
windowsCoder,
183+
windows,
184+
PaneInfo.NO_FIRING,
185+
BeamFnApi.Elements.ElementMetadata.newBuilder().build()))
180186
.setData(dataOutput.toByteString())
181187
.setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(timestamp));
182188
}

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsReshuffleDoFnTest.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import java.util.Arrays;
2525
import java.util.Collection;
2626
import java.util.List;
27+
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
2728
import org.apache.beam.runners.core.DoFnRunner;
2829
import org.apache.beam.runners.core.KeyedWorkItem;
2930
import org.apache.beam.runners.core.NullSideInputReader;
@@ -114,7 +115,12 @@ private <V> void addElement(
114115
valueCoder.encode(value, dataOutput, Context.OUTER);
115116
messageBundle
116117
.addMessagesBuilder()
117-
.setMetadata(WindmillSink.encodeMetadata(windowsCoder, windows, PaneInfo.NO_FIRING))
118+
.setMetadata(
119+
WindmillSink.encodeMetadata(
120+
windowsCoder,
121+
windows,
122+
PaneInfo.NO_FIRING,
123+
BeamFnApi.Elements.ElementMetadata.newBuilder().build()))
118124
.setData(dataOutput.toByteString())
119125
.setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(timestamp));
120126
}

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItemTest.java

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import java.io.IOException;
2323
import java.util.Collection;
2424
import java.util.Collections;
25+
import org.apache.beam.model.fnexecution.v1.BeamFnApi;
2526
import org.apache.beam.runners.core.KeyedWorkItem;
2627
import org.apache.beam.runners.core.StateNamespace;
2728
import org.apache.beam.runners.core.StateNamespaces;
@@ -107,10 +108,14 @@ private void addElement(
107108
long timestamp,
108109
String value,
109110
IntervalWindow window,
110-
PaneInfo paneInfo)
111+
PaneInfo pane)
111112
throws IOException {
112113
ByteString encodedMetadata =
113-
WindmillSink.encodeMetadata(WINDOWS_CODER, Collections.singletonList(window), paneInfo);
114+
WindmillSink.encodeMetadata(
115+
WINDOWS_CODER,
116+
Collections.singletonList(window),
117+
pane,
118+
BeamFnApi.Elements.ElementMetadata.newBuilder().build());
114119
chunk
115120
.addMessagesBuilder()
116121
.setTimestamp(WindmillTimeUtils.harnessToWindmillTimestamp(new Instant(timestamp)))

sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Redistribute.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,7 @@ public Duration getAllowedTimestampSkew() {
179179
public void processElement(
180180
@Element KV<K, ValueInSingleWindow<V>> kv,
181181
OutputReceiver<KV<K, V>> outputReceiver) {
182+
// todo #33176 specify additional metadata in the future
182183
outputReceiver
183184
.builder(KV.of(kv.getKey(), kv.getValue().getValue()))
184185
.setTimestamp(kv.getValue().getTimestamp())

sdks/java/core/src/main/java/org/apache/beam/sdk/transforms/Reify.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,7 @@ public PCollection<KV<K, ValueInSingleWindow<V>>> expand(PCollection<KV<K, V>> i
136136
KvCoder<K, V> coder = (KvCoder<K, V>) input.getCoder();
137137
return input
138138
.apply(
139+
// todo #33176 specify additional metadata in the future
139140
ParDo.of(
140141
new DoFn<KV<K, V>, KV<K, ValueInSingleWindow<V>>>() {
141142
@ProcessElement

0 commit comments

Comments
 (0)