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 3AE84200BA0 for ; Fri, 14 Oct 2016 12:53:18 +0200 (CEST) Received: by cust-asf.ponee.io (Postfix) id 37969160ADD; Fri, 14 Oct 2016 10:53:18 +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 7F487160AD9 for ; Fri, 14 Oct 2016 12:53:17 +0200 (CEST) Received: (qmail 76806 invoked by uid 500); 14 Oct 2016 10:53:16 -0000 Mailing-List: contact commits-help@beam.incubator.apache.org; run by ezmlm Precedence: bulk List-Help: List-Unsubscribe: List-Post: List-Id: Reply-To: dev@beam.incubator.apache.org Delivered-To: mailing list commits@beam.incubator.apache.org Received: (qmail 76797 invoked by uid 99); 14 Oct 2016 10:53:16 -0000 Received: from pnap-us-west-generic-nat.apache.org (HELO spamd2-us-west.apache.org) (209.188.14.142) by apache.org (qpsmtpd/0.29) with ESMTP; Fri, 14 Oct 2016 10:53:16 +0000 Received: from localhost (localhost [127.0.0.1]) by spamd2-us-west.apache.org (ASF Mail Server at spamd2-us-west.apache.org) with ESMTP id 510141A06A2 for ; Fri, 14 Oct 2016 10:53:16 +0000 (UTC) X-Virus-Scanned: Debian amavisd-new at spamd2-us-west.apache.org X-Spam-Flag: NO X-Spam-Score: -6.219 X-Spam-Level: X-Spam-Status: No, score=-6.219 tagged_above=-999 required=6.31 tests=[KAM_ASCII_DIVIDERS=0.8, KAM_LAZY_DOMAIN_SECURITY=1, RCVD_IN_DNSWL_HI=-5, RCVD_IN_MSPIKE_H3=-0.01, RCVD_IN_MSPIKE_WL=-0.01, RP_MATCHES_RCVD=-2.999] autolearn=disabled Received: from mx1-lw-eu.apache.org ([10.40.0.8]) by localhost (spamd2-us-west.apache.org [10.40.0.9]) (amavisd-new, port 10024) with ESMTP id cVdR2vh_eL_V for ; Fri, 14 Oct 2016 10:53:14 +0000 (UTC) Received: from mail.apache.org (hermes.apache.org [140.211.11.3]) by mx1-lw-eu.apache.org (ASF Mail Server at mx1-lw-eu.apache.org) with SMTP id 343C45FB16 for ; Fri, 14 Oct 2016 10:53:13 +0000 (UTC) Received: (qmail 76638 invoked by uid 99); 14 Oct 2016 10:53:12 -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; Fri, 14 Oct 2016 10:53:12 +0000 Received: by git1-us-west.apache.org (ASF Mail Server at git1-us-west.apache.org, from userid 33) id F321CE0381; Fri, 14 Oct 2016 10:53:11 +0000 (UTC) Content-Type: text/plain; charset="us-ascii" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit From: amitsela@apache.org To: commits@beam.incubator.apache.org Date: Fri, 14 Oct 2016 10:53:11 -0000 Message-Id: <6068a1e6a4084d1289dce0cebdf8dfb2@git.apache.org> X-Mailer: ASF-Git Admin Mailer Subject: [1/2] incubator-beam git commit: [BEAM-734] Add support for Spark Streaming Listeners. archived-at: Fri, 14 Oct 2016 10:53:18 -0000 Repository: incubator-beam Updated Branches: refs/heads/master a0f649eac -> d790dfe1b [BEAM-734] Add support for Spark Streaming Listeners. Project: http://git-wip-us.apache.org/repos/asf/incubator-beam/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-beam/commit/4b49abcf Tree: http://git-wip-us.apache.org/repos/asf/incubator-beam/tree/4b49abcf Diff: http://git-wip-us.apache.org/repos/asf/incubator-beam/diff/4b49abcf Branch: refs/heads/master Commit: 4b49abcf7d248e033b2bd8435dff031261f35b73 Parents: a0f649e Author: Sela Authored: Sun Oct 9 13:44:58 2016 +0300 Committer: Sela Committed: Fri Oct 14 12:58:47 2016 +0300 ---------------------------------------------------------------------- .../beam/runners/spark/SparkPipelineOptions.java | 18 ++++++++++++++++++ .../SparkRunnerStreamingContextFactory.java | 8 ++++++++ 2 files changed, 26 insertions(+) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/4b49abcf/runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java ---------------------------------------------------------------------- diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java index 7afb68c..4c20b10 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java @@ -19,6 +19,8 @@ package org.apache.beam.runners.spark; import com.fasterxml.jackson.annotation.JsonIgnore; +import java.util.ArrayList; +import java.util.List; import org.apache.beam.sdk.options.ApplicationNameOptions; import org.apache.beam.sdk.options.Default; import org.apache.beam.sdk.options.DefaultValueFactory; @@ -26,6 +28,8 @@ import org.apache.beam.sdk.options.Description; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.StreamingOptions; import org.apache.spark.api.java.JavaSparkContext; +import org.apache.spark.streaming.api.java.JavaStreamingListener; + /** * Spark runner pipeline options. @@ -88,4 +92,18 @@ public interface SparkPipelineOptions extends PipelineOptions, StreamingOptions, @JsonIgnore JavaSparkContext getProvidedSparkContext(); void setProvidedSparkContext(JavaSparkContext jsc); + + @Description("Spark streaming listeners") + @Default.InstanceFactory(EmptyListenersList.class) + @JsonIgnore + List getListeners(); + void setListeners(List listeners); + + /** Returns an empty list, top avoid handling null. */ + static class EmptyListenersList implements DefaultValueFactory> { + @Override + public List create(PipelineOptions options) { + return new ArrayList<>(); + } + } } http://git-wip-us.apache.org/repos/asf/incubator-beam/blob/4b49abcf/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java ---------------------------------------------------------------------- diff --git a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java index b7a407c..79c87fb 100644 --- a/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java +++ b/runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java @@ -33,6 +33,8 @@ import org.apache.spark.api.java.JavaSparkContext; import org.apache.spark.streaming.Duration; import org.apache.spark.streaming.api.java.JavaStreamingContext; import org.apache.spark.streaming.api.java.JavaStreamingContextFactory; +import org.apache.spark.streaming.api.java.JavaStreamingListener; +import org.apache.spark.streaming.api.java.JavaStreamingListenerWrapper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -89,6 +91,12 @@ public class SparkRunnerStreamingContextFactory implements JavaStreamingContextF } jssc.checkpoint(checkpointDir); + // register listeners. + for (JavaStreamingListener listener: options.getListeners()) { + LOG.info("Registered listener {}." + listener.getClass().getSimpleName()); + jssc.addStreamingListener(new JavaStreamingListenerWrapper(listener)); + } + return jssc; }