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 243D8200C03 for ; Sat, 7 Jan 2017 03:00:57 +0100 (CET) Received: by cust-asf.ponee.io (Postfix) id 22F36160B4C; Sat, 7 Jan 2017 02:00:57 +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 7225C160B39 for ; Sat, 7 Jan 2017 03:00:56 +0100 (CET) Received: (qmail 35522 invoked by uid 500); 7 Jan 2017 02:00:55 -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 35454 invoked by uid 99); 7 Jan 2017 02:00:55 -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; Sat, 07 Jan 2017 02:00:55 +0000 Received: by git1-us-west.apache.org (ASF Mail Server at git1-us-west.apache.org, from userid 33) id 755BADFB6E; Sat, 7 Jan 2017 02:00:55 +0000 (UTC) Content-Type: text/plain; charset="us-ascii" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit From: kenn@apache.org To: commits@beam.apache.org Date: Sat, 07 Jan 2017 02:00:56 -0000 Message-Id: <0334b0f1b821407f9d58f04a60b858ba@git.apache.org> In-Reply-To: References: X-Mailer: ASF-Git Admin Mailer Subject: [2/3] beam git commit: Rollforwards "Allow stateful DoFn in DataflowRunner"" archived-at: Sat, 07 Jan 2017 02:00:57 -0000 Rollforwards "Allow stateful DoFn in DataflowRunner"" This rolls forward 42bb15d, allowing stateful DoFn in DataflowRunner. Project: http://git-wip-us.apache.org/repos/asf/beam/repo Commit: http://git-wip-us.apache.org/repos/asf/beam/commit/d583a1ce Tree: http://git-wip-us.apache.org/repos/asf/beam/tree/d583a1ce Diff: http://git-wip-us.apache.org/repos/asf/beam/diff/d583a1ce Branch: refs/heads/master Commit: d583a1cedfd0c6abc9ca5059009965186b51e040 Parents: 5af7f42 Author: Kenneth Knowles Authored: Thu Jan 5 19:15:24 2017 -0800 Committer: Kenneth Knowles Committed: Fri Jan 6 14:04:05 2017 -0800 ---------------------------------------------------------------------- runners/google-cloud-dataflow-java/pom.xml | 1 - .../dataflow/DataflowPipelineTranslator.java | 22 +++++++------------- 2 files changed, 8 insertions(+), 15 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/beam/blob/d583a1ce/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 64727b1..7bf2089 100644 --- a/runners/google-cloud-dataflow-java/pom.xml +++ b/runners/google-cloud-dataflow-java/pom.xml @@ -78,7 +78,6 @@ runnable-on-service-tests - org.apache.beam.sdk.testing.UsesStatefulParDo, org.apache.beam.sdk.testing.UsesTimersInParDo, org.apache.beam.sdk.testing.UsesSplittableParDo, org.apache.beam.sdk.testing.UsesMetrics http://git-wip-us.apache.org/repos/asf/beam/blob/d583a1ce/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 8e5901e..524c30b 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 @@ -79,6 +79,7 @@ import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.transforms.View; import org.apache.beam.sdk.transforms.display.DisplayData; import org.apache.beam.sdk.transforms.display.HasDisplayData; +import org.apache.beam.sdk.transforms.reflect.DoFnSignature; import org.apache.beam.sdk.transforms.reflect.DoFnSignatures; import org.apache.beam.sdk.transforms.windowing.DefaultTrigger; import org.apache.beam.sdk.transforms.windowing.Window; @@ -836,7 +837,6 @@ class DataflowPipelineTranslator { private void translateMultiHelper( ParDo.BoundMulti transform, TranslationContext context) { - DataflowPipelineTranslator.rejectStatefulDoFn(transform.getFn()); StepTranslationContext stepContext = context.addStep(transform, "ParallelDo"); translateInputs( @@ -865,7 +865,6 @@ class DataflowPipelineTranslator { private void translateSingleHelper( ParDo.Bound transform, TranslationContext context) { - rejectStatefulDoFn(transform.getFn()); StepTranslationContext stepContext = context.addStep(transform, "ParallelDo"); translateInputs( @@ -914,18 +913,6 @@ class DataflowPipelineTranslator { registerTransformTranslator(Read.Bounded.class, new ReadTranslator()); } - private static void rejectStatefulDoFn(DoFn doFn) { - if (DoFnSignatures.getSignature(doFn.getClass()).isStateful()) { - throw new UnsupportedOperationException( - String.format( - "Found %s annotations on %s, but %s cannot yet be used with state in the %s.", - DoFn.StateId.class.getSimpleName(), - doFn.getClass().getName(), - DoFn.class.getSimpleName(), - DataflowRunner.class.getSimpleName())); - } - } - private static void translateInputs( StepTranslationContext stepContext, PCollection input, @@ -960,6 +947,9 @@ class DataflowPipelineTranslator { TranslationContext context, long mainOutput, Map> outputMap) { + + DoFnSignature signature = DoFnSignatures.getSignature(fn.getClass()); + stepContext.addInput(PropertyNames.USER_FN, fn.getClass().getName()); stepContext.addInput( PropertyNames.SERIALIZED_FN, @@ -967,6 +957,10 @@ class DataflowPipelineTranslator { serializeToByteArray( DoFnInfo.forFn( fn, windowingStrategy, sideInputs, inputCoder, mainOutput, outputMap)))); + + if (signature.isStateful()) { + stepContext.addInput(PropertyNames.USES_KEYED_STATE, "true"); + } } private static BiMap> translateOutputs(