scwhittle commented on code in PR #39961:
URL: https://github.com/apache/beam/pull/39961#discussion_r3947927950
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubSink.java:
##########
@@ -187,18 +190,41 @@ public long add(WindowedValue<T> data) throws IOException
{
return byteString.size();
}
+ private void flush(boolean bundleLevel) {
+ try {
+ Windmill.PubSubMessageBundle pubsubMessages = outputBuilder.build();
+ if (pubsubMessages.getMessagesCount() > 0) {
+ if (bundleLevel) {
+ context.addBundlePubsubMessages(pubsubMessages);
+ } else {
+ context.getOutputBuilder().addPubsubMessages(pubsubMessages);
+ }
+ }
+ } finally {
+ outputBuilder = createOutputBuilder();
+ }
+ }
+
+ @Override
+ public void finishKey(@Nullable Object key) throws IOException {
+ if (context.multiKeyBundleEnabled()) {
+ flush(/* bundleLevel= */ false);
+ }
+ }
+
@Override
public void close() throws IOException {
- Windmill.PubSubMessageBundle pubsubMessages = outputBuilder.build();
- if (pubsubMessages.getMessagesCount() > 0) {
- context.getOutputBuilder().addPubsubMessages(pubsubMessages);
+ if (context.multiKeyBundleEnabled()) {
Review Comment:
just pass multiKeyBundleEnabled() as param?
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubDynamicSink.java:
##########
@@ -145,14 +139,43 @@ public long add(WindowedValue<PubsubMessage> data) throws
IOException {
return byteString.size();
}
+ private void flush(boolean bundleLevel) {
+ try {
+ for (Windmill.PubSubMessageBundle.Builder builder :
outputBuilders.values()) {
+ if (builder.getMessagesCount() > 0) {
+ Windmill.PubSubMessageBundle pubsubMessages = builder.build();
+ if (bundleLevel) {
+ context.addBundlePubsubMessages(pubsubMessages);
+ } else {
+ context.getOutputBuilder().addPubsubMessages(pubsubMessages);
+ }
+ }
+ }
+ } finally {
+ outputBuilders.clear();
+ }
+ }
+
+ @Override
+ public void finishKey(@Nullable Object key) throws IOException {
+ if (context.multiKeyBundleEnabled()) {
+ flush(/* bundleLevel= */ false);
+ }
+ }
+
@Override
public void close() throws IOException {
- outputBuilders.clear();
+ if (context.multiKeyBundleEnabled()) {
Review Comment:
just pass multiKeyBundleEnabled() as param to flush
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubSink.java:
##########
@@ -187,18 +190,41 @@ public long add(WindowedValue<T> data) throws IOException
{
return byteString.size();
}
+ private void flush(boolean bundleLevel) {
+ try {
+ Windmill.PubSubMessageBundle pubsubMessages = outputBuilder.build();
+ if (pubsubMessages.getMessagesCount() > 0) {
+ if (bundleLevel) {
+ context.addBundlePubsubMessages(pubsubMessages);
Review Comment:
maybe we should try to merge with existing bundles? If we have 100 keys each
publishing 1 message we're going to end up with 100 separate bundles and
duplicating the topic etc, instead of 1 bundle.
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubDynamicSink.java:
##########
@@ -145,14 +139,43 @@ public long add(WindowedValue<PubsubMessage> data) throws
IOException {
return byteString.size();
}
+ private void flush(boolean bundleLevel) {
+ try {
+ for (Windmill.PubSubMessageBundle.Builder builder :
outputBuilders.values()) {
+ if (builder.getMessagesCount() > 0) {
+ Windmill.PubSubMessageBundle pubsubMessages = builder.build();
+ if (bundleLevel) {
+ context.addBundlePubsubMessages(pubsubMessages);
Review Comment:
see comments in PubsubSink about merging etc
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubSink.java:
##########
@@ -187,18 +190,41 @@ public long add(WindowedValue<T> data) throws IOException
{
return byteString.size();
}
+ private void flush(boolean bundleLevel) {
+ try {
+ Windmill.PubSubMessageBundle pubsubMessages = outputBuilder.build();
+ if (pubsubMessages.getMessagesCount() > 0) {
+ if (bundleLevel) {
+ context.addBundlePubsubMessages(pubsubMessages);
+ } else {
+ context.getOutputBuilder().addPubsubMessages(pubsubMessages);
+ }
+ }
+ } finally {
+ outputBuilder = createOutputBuilder();
+ }
+ }
+
+ @Override
+ public void finishKey(@Nullable Object key) throws IOException {
+ if (context.multiKeyBundleEnabled()) {
+ flush(/* bundleLevel= */ false);
Review Comment:
see above comment, is there a benefit to producing at the key level and not
the bundle level? With bundle level we can reduce duplicating topic info and
proto overhead
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/Sink.java:
##########
@@ -36,6 +37,9 @@ public interface SinkWriter<ElemT> extends AutoCloseable {
/** Adds a value to the sink. Returns the size in bytes of the data
written. */
public long add(ElemT value) throws IOException;
+ /** Called when all elements for a specific key have been processed. */
Review Comment:
Should we note this is optional and only currently done for multi-key
streaming execution?
##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java:
##########
@@ -350,26 +350,47 @@ public long add(WindowedValue<T> data) throws IOException
{
return (long) key.size() + value.size() + metadata.size() + id.size() +
offsetSize;
}
- @Override
- public void close() throws IOException {
+ private void flush(boolean bundleLevel) {
try {
outputBuilder.setDestinationStreamId(destinationName);
for (Windmill.KeyedMessageBundle.Builder keyedOutput :
productionMap.values()) {
outputBuilder.addBundles(keyedOutput.build());
}
if (outputBuilder.getBundlesCount() > 0) {
- context.getOutputBuilder().addOutputMessages(outputBuilder.build());
+ Windmill.OutputMessageBundle bundle = outputBuilder.build();
+ if (bundleLevel) {
+ context.addBundleOutputMessages(bundle);
Review Comment:
same merging comments
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]