Return-Path: X-Original-To: archive-asf-public-internal@cust-asf2.ponee.io Delivered-To: archive-asf-public-internal@cust-asf2.ponee.io Received: from cust-asf.ponee.io (cust-asf.ponee.io [163.172.22.183]) by cust-asf2.ponee.io (Postfix) with ESMTP id 2B812200C54 for ; Wed, 12 Apr 2017 21:56:36 +0200 (CEST) Received: by cust-asf.ponee.io (Postfix) id 2A2B9160BB5; Wed, 12 Apr 2017 19:56:36 +0000 (UTC) Delivered-To: archive-asf-public@cust-asf.ponee.io Received: from mail.apache.org (hermes.apache.org [140.211.11.3]) by cust-asf.ponee.io (Postfix) with SMTP id 2533C160BAF for ; Wed, 12 Apr 2017 21:56:34 +0200 (CEST) Received: (qmail 61452 invoked by uid 500); 12 Apr 2017 19:56:33 -0000 Mailing-List: contact commits-help@beam.apache.org; run by ezmlm Precedence: bulk List-Help: List-Unsubscribe: List-Post: List-Id: Reply-To: dev@beam.apache.org Delivered-To: mailing list commits@beam.apache.org Received: (qmail 60296 invoked by uid 99); 12 Apr 2017 19:56:33 -0000 Received: from git1-us-west.apache.org (HELO git1-us-west.apache.org) (140.211.11.23) by apache.org (qpsmtpd/0.29) with ESMTP; Wed, 12 Apr 2017 19:56:33 +0000 Received: by git1-us-west.apache.org (ASF Mail Server at git1-us-west.apache.org, from userid 33) id 1243EDFC00; Wed, 12 Apr 2017 19:56:32 +0000 (UTC) Content-Type: text/plain; charset="us-ascii" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit From: jbonofre@apache.org To: commits@beam.apache.org Date: Wed, 12 Apr 2017 19:56:49 -0000 Message-Id: <0b7003dc095c46d29f660c34d538fd07@git.apache.org> In-Reply-To: <75c8ab1209c4423a85fad7497d6c2c70@git.apache.org> References: <75c8ab1209c4423a85fad7497d6c2c70@git.apache.org> X-Mailer: ASF-Git Admin Mailer Subject: [18/50] [abbrv] beam git commit: DataflowRunner: send windowing strategy using Runner API proto archived-at: Wed, 12 Apr 2017 19:56:36 -0000 DataflowRunner: send windowing strategy using Runner API proto Project: http://git-wip-us.apache.org/repos/asf/beam/repo Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/69f412dc Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/69f412dc Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/69f412dc Branch: refs/heads/DSL_SQL Commit: 69f412dc34f36df40b034c2160b8b0cdad815011 Parents: bef2d37 Author: Dan Halperin Authored: Tue Apr 11 13:32:11 2017 -0700 Committer: Dan Halperin Committed: Tue Apr 11 15:04:38 2017 -0700 ---------------------------------------------------------------------- runners/google-cloud-dataflow-java/pom.xml | 12 +++++++++++- .../runners/dataflow/DataflowPipelineTranslator.java | 14 ++++++++++++-- 2 files changed, 23 insertions(+), 3 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/beam/blob/69f412dc/runners/google-cloud-dataflow-java/pom.xml ---------------------------------------------------------------------- diff --git a/runners/google-cloud-dataflow-java/pom.xml b/runners/google-cloud-dataflow-java/pom.xml index d0d86e6..a57744c 100644 --- a/runners/google-cloud-dataflow-java/pom.xml +++ b/runners/google-cloud-dataflow-java/pom.xml @@ -33,7 +33,7 @@ jar - beam-master-20170410 + beam-master-20170411 1 6 @@ -114,6 +114,7 @@ com.google.cloud.bigtable:bigtable-client-core com.google.guava:guava + org.apache.beam:beam-runners-core-construction-java @@ -153,6 +154,10 @@ com.google.cloud.bigtable.grpc.BigtableTableName + + org.apache.beam.runners.core + org.apache.beam.runners.dataflow.repackaged.org.apache.beam.runners.core + @@ -178,6 +183,11 @@ org.apache.beam + beam-sdks-common-runner-api + + + + org.apache.beam beam-runners-core-construction-java http://git-wip-us.apache.org/repos/asf/beam/blob/69f412dc/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java ---------------------------------------------------------------------- diff --git a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java index 34da996..abeca4d 100644 --- a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java +++ b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java @@ -55,6 +55,7 @@ import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicLong; import javax.annotation.Nullable; +import org.apache.beam.runners.core.construction.WindowingStrategies; import org.apache.beam.runners.dataflow.BatchViewOverrides.GroupByKeyAndSortValuesOnly; import org.apache.beam.runners.dataflow.DataflowRunner.CombineGroupedValues; import org.apache.beam.runners.dataflow.PrimitiveParDoSingleFactory.ParDoSingle; @@ -111,6 +112,15 @@ public class DataflowPipelineTranslator { private static final Logger LOG = LoggerFactory.getLogger(DataflowPipelineTranslator.class); private static final ObjectMapper MAPPER = new ObjectMapper(); + private static byte[] serializeWindowingStrategy(WindowingStrategy windowingStrategy) { + try { + return WindowingStrategies.toProto(windowingStrategy).toByteArray(); + } catch (Exception e) { + throw new RuntimeException( + String.format("Unable to format windowing strategy %s as bytes", windowingStrategy), e); + } + } + /** * A map from {@link PTransform} subclass to the corresponding * {@link TransformTranslator} to use to translate that transform. @@ -813,7 +823,7 @@ public class DataflowPipelineTranslator { stepContext.addInput(PropertyNames.DISALLOW_COMBINER_LIFTING, disallowCombinerLifting); stepContext.addInput( PropertyNames.SERIALIZED_FN, - byteArrayToJsonString(serializeToByteArray(windowingStrategy))); + byteArrayToJsonString(serializeWindowingStrategy(windowingStrategy))); stepContext.addInput( PropertyNames.IS_MERGING_WINDOW_FN, !windowingStrategy.getWindowFn().isNonMerging()); @@ -891,7 +901,7 @@ public class DataflowPipelineTranslator { stepContext.addOutput(context.getOutput(transform)); WindowingStrategy strategy = context.getOutput(transform).getWindowingStrategy(); - byte[] serializedBytes = serializeToByteArray(strategy); + byte[] serializedBytes = serializeWindowingStrategy(strategy); String serializedJson = byteArrayToJsonString(serializedBytes); stepContext.addInput(PropertyNames.SERIALIZED_FN, serializedJson); }