This is an automated email from the ASF dual-hosted git repository. echauchot pushed a commit to branch spark-runner_structured-streaming in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/spark-runner_structured-streaming by this push: new 0402d6a Cleaning 0402d6a is described below commit 0402d6a0a379b973fba68524ccaf6ab2ea061d2c Author: Etienne Chauchot <echauc...@apache.org> AuthorDate: Tue Jan 15 17:39:27 2019 +0100 Cleaning --- .../org/apache/beam/runners/spark/structuredstreaming/SparkRunner.java | 1 + .../spark/structuredstreaming/translation/TranslationContext.java | 3 +-- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkRunner.java b/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkRunner.java index 934c6d2..72cb524 100644 --- a/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkRunner.java +++ b/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkRunner.java @@ -114,6 +114,7 @@ public final class SparkRunner extends PipelineRunner<SparkPipelineResult> { public SparkPipelineResult run(final Pipeline pipeline) { translationContext = translatePipeline(pipeline); //TODO initialise other services: checkpointing, metrics system, listeners, ... + //TODO pass testMode using pipelineOptions translationContext.startPipeline(true); return new SparkPipelineResult(); } diff --git a/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TranslationContext.java b/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TranslationContext.java index 0f20663..75b470e 100644 --- a/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TranslationContext.java +++ b/runners/spark-structured-streaming/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TranslationContext.java @@ -188,8 +188,7 @@ public class TranslationContext { } } else { // apply a dummy fn just to apply forech action that will trigger the pipeline run in spark - dataset.foreachPartition(t -> { - }); + dataset.foreachPartition(t -> {}); } } }