<!DOCTYPE html PUBLIC "-//W3C//DTD XHTML 1.0 Transitional//EN"
"http://www.w3.org/TR/xhtml1/DTD/xhtml1-transitional.dtd">
<html xmlns="http://www.w3.org/1999/xhtml">
<head>
<meta http-equiv="Content-Type" content="text/html; charset=utf-8" />
<title>pyspark.streaming.kafka — PySpark 1.5.0 documentation</title>
<link rel="stylesheet" href="../../../_static/nature.css" type="text/css" />
<link rel="stylesheet" href="../../../_static/pygments.css" type="text/css" />
<script type="text/javascript">
var DOCUMENTATION_OPTIONS = {
URL_ROOT: '../../../',
VERSION: '1.5.0',
COLLAPSE_INDEX: false,
FILE_SUFFIX: '.html',
HAS_SOURCE: true
};
</script>
<script type="text/javascript" src="../../../_static/jquery.js"></script>
<script type="text/javascript" src="../../../_static/underscore.js"></script>
<script type="text/javascript" src="../../../_static/doctools.js"></script>
<link rel="top" title="PySpark 1.5.0 documentation" href="../../../index.html" />
<link rel="up" title="Module code" href="../../index.html" />
</head>
<body role="document">
<div class="related" role="navigation" aria-label="related navigation">
<h3>Navigation</h3>
<ul>
<li class="nav-item nav-item-0"><a href="../../../index.html">PySpark 1.5.0 documentation</a> »</li>
<li class="nav-item nav-item-1"><a href="../../index.html" accesskey="U">Module code</a> »</li>
</ul>
</div>
<div class="document">
<div class="documentwrapper">
<div class="bodywrapper">
<div class="body" role="main">
<h1>Source code for pyspark.streaming.kafka</h1><div class="highlight"><pre>
<span class="c">#</span>
<span class="c"># Licensed to the Apache Software Foundation (ASF) under one or more</span>
<span class="c"># contributor license agreements. See the NOTICE file distributed with</span>
<span class="c"># this work for additional information regarding copyright ownership.</span>
<span class="c"># The ASF licenses this file to You under the Apache License, Version 2.0</span>
<span class="c"># (the "License"); you may not use this file except in compliance with</span>
<span class="c"># the License. You may obtain a copy of the License at</span>
<span class="c">#</span>
<span class="c"># http://www.apache.org/licenses/LICENSE-2.0</span>
<span class="c">#</span>
<span class="c"># Unless required by applicable law or agreed to in writing, software</span>
<span class="c"># distributed under the License is distributed on an "AS IS" BASIS,</span>
<span class="c"># WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.</span>
<span class="c"># See the License for the specific language governing permissions and</span>
<span class="c"># limitations under the License.</span>
<span class="c">#</span>
<span class="kn">from</span> <span class="nn">py4j.java_gateway</span> <span class="kn">import</span> <span class="n">Py4JJavaError</span>
<span class="kn">from</span> <span class="nn">pyspark.rdd</span> <span class="kn">import</span> <span class="n">RDD</span>
<span class="kn">from</span> <span class="nn">pyspark.storagelevel</span> <span class="kn">import</span> <span class="n">StorageLevel</span>
<span class="kn">from</span> <span class="nn">pyspark.serializers</span> <span class="kn">import</span> <span class="n">PairDeserializer</span><span class="p">,</span> <span class="n">NoOpSerializer</span>
<span class="kn">from</span> <span class="nn">pyspark.streaming</span> <span class="kn">import</span> <span class="n">DStream</span>
<span class="kn">from</span> <span class="nn">pyspark.streaming.dstream</span> <span class="kn">import</span> <span class="n">TransformedDStream</span>
<span class="kn">from</span> <span class="nn">pyspark.streaming.util</span> <span class="kn">import</span> <span class="n">TransformFunction</span>
<span class="n">__all__</span> <span class="o">=</span> <span class="p">[</span><span class="s">'Broker'</span><span class="p">,</span> <span class="s">'KafkaUtils'</span><span class="p">,</span> <span class="s">'OffsetRange'</span><span class="p">,</span> <span class="s">'TopicAndPartition'</span><span class="p">,</span> <span class="s">'utf8_decoder'</span><span class="p">]</span>
<div class="viewcode-block" id="utf8_decoder"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.utf8_decoder">[docs]</a><span class="k">def</span> <span class="nf">utf8_decoder</span><span class="p">(</span><span class="n">s</span><span class="p">):</span>
<span class="sd">""" Decode the unicode as UTF-8 """</span>
<span class="k">if</span> <span class="n">s</span> <span class="ow">is</span> <span class="bp">None</span><span class="p">:</span>
<span class="k">return</span> <span class="bp">None</span>
<span class="k">return</span> <span class="n">s</span><span class="o">.</span><span class="n">decode</span><span class="p">(</span><span class="s">'utf-8'</span><span class="p">)</span>
</div>
<div class="viewcode-block" id="KafkaUtils"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.KafkaUtils">[docs]</a><span class="k">class</span> <span class="nc">KafkaUtils</span><span class="p">(</span><span class="nb">object</span><span class="p">):</span>
<span class="nd">@staticmethod</span>
<div class="viewcode-block" id="KafkaUtils.createStream"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.KafkaUtils.createStream">[docs]</a> <span class="k">def</span> <span class="nf">createStream</span><span class="p">(</span><span class="n">ssc</span><span class="p">,</span> <span class="n">zkQuorum</span><span class="p">,</span> <span class="n">groupId</span><span class="p">,</span> <span class="n">topics</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="o">=</span><span class="bp">None</span><span class="p">,</span>
<span class="n">storageLevel</span><span class="o">=</span><span class="n">StorageLevel</span><span class="o">.</span><span class="n">MEMORY_AND_DISK_SER_2</span><span class="p">,</span>
<span class="n">keyDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">,</span> <span class="n">valueDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Create an input stream that pulls messages from a Kafka Broker.</span>
<span class="sd"> :param ssc: StreamingContext object</span>
<span class="sd"> :param zkQuorum: Zookeeper quorum (hostname:port,hostname:port,..).</span>
<span class="sd"> :param groupId: The group id for this consumer.</span>
<span class="sd"> :param topics: Dict of (topic_name -> numPartitions) to consume.</span>
<span class="sd"> Each partition is consumed in its own thread.</span>
<span class="sd"> :param kafkaParams: Additional params for Kafka</span>
<span class="sd"> :param storageLevel: RDD storage level.</span>
<span class="sd"> :param keyDecoder: A function used to decode key (default is utf8_decoder)</span>
<span class="sd"> :param valueDecoder: A function used to decode value (default is utf8_decoder)</span>
<span class="sd"> :return: A DStream object</span>
<span class="sd"> """</span>
<span class="k">if</span> <span class="n">kafkaParams</span> <span class="ow">is</span> <span class="bp">None</span><span class="p">:</span>
<span class="n">kafkaParams</span> <span class="o">=</span> <span class="nb">dict</span><span class="p">()</span>
<span class="n">kafkaParams</span><span class="o">.</span><span class="n">update</span><span class="p">({</span>
<span class="s">"zookeeper.connect"</span><span class="p">:</span> <span class="n">zkQuorum</span><span class="p">,</span>
<span class="s">"group.id"</span><span class="p">:</span> <span class="n">groupId</span><span class="p">,</span>
<span class="s">"zookeeper.connection.timeout.ms"</span><span class="p">:</span> <span class="s">"10000"</span><span class="p">,</span>
<span class="p">})</span>
<span class="k">if</span> <span class="ow">not</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">topics</span><span class="p">,</span> <span class="nb">dict</span><span class="p">):</span>
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s">"topics should be dict"</span><span class="p">)</span>
<span class="n">jlevel</span> <span class="o">=</span> <span class="n">ssc</span><span class="o">.</span><span class="n">_sc</span><span class="o">.</span><span class="n">_getJavaStorageLevel</span><span class="p">(</span><span class="n">storageLevel</span><span class="p">)</span>
<span class="k">try</span><span class="p">:</span>
<span class="c"># Use KafkaUtilsPythonHelper to access Scala's KafkaUtils (see SPARK-6027)</span>
<span class="n">helperClass</span> <span class="o">=</span> <span class="n">ssc</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">java</span><span class="o">.</span><span class="n">lang</span><span class="o">.</span><span class="n">Thread</span><span class="o">.</span><span class="n">currentThread</span><span class="p">()</span><span class="o">.</span><span class="n">getContextClassLoader</span><span class="p">()</span>\
<span class="o">.</span><span class="n">loadClass</span><span class="p">(</span><span class="s">"org.apache.spark.streaming.kafka.KafkaUtilsPythonHelper"</span><span class="p">)</span>
<span class="n">helper</span> <span class="o">=</span> <span class="n">helperClass</span><span class="o">.</span><span class="n">newInstance</span><span class="p">()</span>
<span class="n">jstream</span> <span class="o">=</span> <span class="n">helper</span><span class="o">.</span><span class="n">createStream</span><span class="p">(</span><span class="n">ssc</span><span class="o">.</span><span class="n">_jssc</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="p">,</span> <span class="n">topics</span><span class="p">,</span> <span class="n">jlevel</span><span class="p">)</span>
<span class="k">except</span> <span class="n">Py4JJavaError</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
<span class="c"># TODO: use --jar once it also work on driver</span>
<span class="k">if</span> <span class="s">'ClassNotFoundException'</span> <span class="ow">in</span> <span class="nb">str</span><span class="p">(</span><span class="n">e</span><span class="o">.</span><span class="n">java_exception</span><span class="p">):</span>
<span class="n">KafkaUtils</span><span class="o">.</span><span class="n">_printErrorMsg</span><span class="p">(</span><span class="n">ssc</span><span class="o">.</span><span class="n">sparkContext</span><span class="p">)</span>
<span class="k">raise</span> <span class="n">e</span>
<span class="n">ser</span> <span class="o">=</span> <span class="n">PairDeserializer</span><span class="p">(</span><span class="n">NoOpSerializer</span><span class="p">(),</span> <span class="n">NoOpSerializer</span><span class="p">())</span>
<span class="n">stream</span> <span class="o">=</span> <span class="n">DStream</span><span class="p">(</span><span class="n">jstream</span><span class="p">,</span> <span class="n">ssc</span><span class="p">,</span> <span class="n">ser</span><span class="p">)</span>
<span class="k">return</span> <span class="n">stream</span><span class="o">.</span><span class="n">map</span><span class="p">(</span><span class="k">lambda</span> <span class="n">k_v</span><span class="p">:</span> <span class="p">(</span><span class="n">keyDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">0</span><span class="p">]),</span> <span class="n">valueDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">1</span><span class="p">])))</span>
</div>
<span class="nd">@staticmethod</span>
<div class="viewcode-block" id="KafkaUtils.createDirectStream"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.KafkaUtils.createDirectStream">[docs]</a> <span class="k">def</span> <span class="nf">createDirectStream</span><span class="p">(</span><span class="n">ssc</span><span class="p">,</span> <span class="n">topics</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="p">,</span> <span class="n">fromOffsets</span><span class="o">=</span><span class="bp">None</span><span class="p">,</span>
<span class="n">keyDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">,</span> <span class="n">valueDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> .. note:: Experimental</span>
<span class="sd"> Create an input stream that directly pulls messages from a Kafka Broker and specific offset.</span>
<span class="sd"> This is not a receiver based Kafka input stream, it directly pulls the message from Kafka</span>
<span class="sd"> in each batch duration and processed without storing.</span>
<span class="sd"> This does not use Zookeeper to store offsets. The consumed offsets are tracked</span>
<span class="sd"> by the stream itself. For interoperability with Kafka monitoring tools that depend on</span>
<span class="sd"> Zookeeper, you have to update Kafka/Zookeeper yourself from the streaming application.</span>
<span class="sd"> You can access the offsets used in each batch from the generated RDDs (see</span>
<span class="sd"> To recover from driver failures, you have to enable checkpointing in the StreamingContext.</span>
<span class="sd"> The information on consumed offset can be recovered from the checkpoint.</span>
<span class="sd"> See the programming guide for details (constraints, etc.).</span>
<span class="sd"> :param ssc: StreamingContext object.</span>
<span class="sd"> :param topics: list of topic_name to consume.</span>
<span class="sd"> :param kafkaParams: Additional params for Kafka.</span>
<span class="sd"> :param fromOffsets: Per-topic/partition Kafka offsets defining the (inclusive) starting</span>
<span class="sd"> point of the stream.</span>
<span class="sd"> :param keyDecoder: A function used to decode key (default is utf8_decoder).</span>
<span class="sd"> :param valueDecoder: A function used to decode value (default is utf8_decoder).</span>
<span class="sd"> :return: A DStream object</span>
<span class="sd"> """</span>
<span class="k">if</span> <span class="n">fromOffsets</span> <span class="ow">is</span> <span class="bp">None</span><span class="p">:</span>
<span class="n">fromOffsets</span> <span class="o">=</span> <span class="nb">dict</span><span class="p">()</span>
<span class="k">if</span> <span class="ow">not</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">topics</span><span class="p">,</span> <span class="nb">list</span><span class="p">):</span>
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s">"topics should be list"</span><span class="p">)</span>
<span class="k">if</span> <span class="ow">not</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">kafkaParams</span><span class="p">,</span> <span class="nb">dict</span><span class="p">):</span>
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s">"kafkaParams should be dict"</span><span class="p">)</span>
<span class="k">try</span><span class="p">:</span>
<span class="n">helperClass</span> <span class="o">=</span> <span class="n">ssc</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">java</span><span class="o">.</span><span class="n">lang</span><span class="o">.</span><span class="n">Thread</span><span class="o">.</span><span class="n">currentThread</span><span class="p">()</span><span class="o">.</span><span class="n">getContextClassLoader</span><span class="p">()</span> \
<span class="o">.</span><span class="n">loadClass</span><span class="p">(</span><span class="s">"org.apache.spark.streaming.kafka.KafkaUtilsPythonHelper"</span><span class="p">)</span>
<span class="n">helper</span> <span class="o">=</span> <span class="n">helperClass</span><span class="o">.</span><span class="n">newInstance</span><span class="p">()</span>
<span class="n">jfromOffsets</span> <span class="o">=</span> <span class="nb">dict</span><span class="p">([(</span><span class="n">k</span><span class="o">.</span><span class="n">_jTopicAndPartition</span><span class="p">(</span><span class="n">helper</span><span class="p">),</span>
<span class="n">v</span><span class="p">)</span> <span class="k">for</span> <span class="p">(</span><span class="n">k</span><span class="p">,</span> <span class="n">v</span><span class="p">)</span> <span class="ow">in</span> <span class="n">fromOffsets</span><span class="o">.</span><span class="n">items</span><span class="p">()])</span>
<span class="n">jstream</span> <span class="o">=</span> <span class="n">helper</span><span class="o">.</span><span class="n">createDirectStream</span><span class="p">(</span><span class="n">ssc</span><span class="o">.</span><span class="n">_jssc</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="p">,</span> <span class="nb">set</span><span class="p">(</span><span class="n">topics</span><span class="p">),</span> <span class="n">jfromOffsets</span><span class="p">)</span>
<span class="k">except</span> <span class="n">Py4JJavaError</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
<span class="k">if</span> <span class="s">'ClassNotFoundException'</span> <span class="ow">in</span> <span class="nb">str</span><span class="p">(</span><span class="n">e</span><span class="o">.</span><span class="n">java_exception</span><span class="p">):</span>
<span class="n">KafkaUtils</span><span class="o">.</span><span class="n">_printErrorMsg</span><span class="p">(</span><span class="n">ssc</span><span class="o">.</span><span class="n">sparkContext</span><span class="p">)</span>
<span class="k">raise</span> <span class="n">e</span>
<span class="n">ser</span> <span class="o">=</span> <span class="n">PairDeserializer</span><span class="p">(</span><span class="n">NoOpSerializer</span><span class="p">(),</span> <span class="n">NoOpSerializer</span><span class="p">())</span>
<span class="n">stream</span> <span class="o">=</span> <span class="n">DStream</span><span class="p">(</span><span class="n">jstream</span><span class="p">,</span> <span class="n">ssc</span><span class="p">,</span> <span class="n">ser</span><span class="p">)</span> \
<span class="o">.</span><span class="n">map</span><span class="p">(</span><span class="k">lambda</span> <span class="n">k_v</span><span class="p">:</span> <span class="p">(</span><span class="n">keyDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">0</span><span class="p">]),</span> <span class="n">valueDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">1</span><span class="p">])))</span>
<span class="k">return</span> <span class="n">KafkaDStream</span><span class="p">(</span><span class="n">stream</span><span class="o">.</span><span class="n">_jdstream</span><span class="p">,</span> <span class="n">ssc</span><span class="p">,</span> <span class="n">stream</span><span class="o">.</span><span class="n">_jrdd_deserializer</span><span class="p">)</span>
</div>
<span class="nd">@staticmethod</span>
<div class="viewcode-block" id="KafkaUtils.createRDD"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.KafkaUtils.createRDD">[docs]</a> <span class="k">def</span> <span class="nf">createRDD</span><span class="p">(</span><span class="n">sc</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="p">,</span> <span class="n">offsetRanges</span><span class="p">,</span> <span class="n">leaders</span><span class="o">=</span><span class="bp">None</span><span class="p">,</span>
<span class="n">keyDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">,</span> <span class="n">valueDecoder</span><span class="o">=</span><span class="n">utf8_decoder</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> .. note:: Experimental</span>
<span class="sd"> Create a RDD from Kafka using offset ranges for each topic and partition.</span>
<span class="sd"> :param sc: SparkContext object</span>
<span class="sd"> :param kafkaParams: Additional params for Kafka</span>
<span class="sd"> :param offsetRanges: list of offsetRange to specify topic:partition:[start, end) to consume</span>
<span class="sd"> :param leaders: Kafka brokers for each TopicAndPartition in offsetRanges. May be an empty</span>
<span class="sd"> map, in which case leaders will be looked up on the driver.</span>
<span class="sd"> :param keyDecoder: A function used to decode key (default is utf8_decoder)</span>
<span class="sd"> :param valueDecoder: A function used to decode value (default is utf8_decoder)</span>
<span class="sd"> :return: A RDD object</span>
<span class="sd"> """</span>
<span class="k">if</span> <span class="n">leaders</span> <span class="ow">is</span> <span class="bp">None</span><span class="p">:</span>
<span class="n">leaders</span> <span class="o">=</span> <span class="nb">dict</span><span class="p">()</span>
<span class="k">if</span> <span class="ow">not</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">kafkaParams</span><span class="p">,</span> <span class="nb">dict</span><span class="p">):</span>
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s">"kafkaParams should be dict"</span><span class="p">)</span>
<span class="k">if</span> <span class="ow">not</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">offsetRanges</span><span class="p">,</span> <span class="nb">list</span><span class="p">):</span>
<span class="k">raise</span> <span class="ne">TypeError</span><span class="p">(</span><span class="s">"offsetRanges should be list"</span><span class="p">)</span>
<span class="k">try</span><span class="p">:</span>
<span class="n">helperClass</span> <span class="o">=</span> <span class="n">sc</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">java</span><span class="o">.</span><span class="n">lang</span><span class="o">.</span><span class="n">Thread</span><span class="o">.</span><span class="n">currentThread</span><span class="p">()</span><span class="o">.</span><span class="n">getContextClassLoader</span><span class="p">()</span> \
<span class="o">.</span><span class="n">loadClass</span><span class="p">(</span><span class="s">"org.apache.spark.streaming.kafka.KafkaUtilsPythonHelper"</span><span class="p">)</span>
<span class="n">helper</span> <span class="o">=</span> <span class="n">helperClass</span><span class="o">.</span><span class="n">newInstance</span><span class="p">()</span>
<span class="n">joffsetRanges</span> <span class="o">=</span> <span class="p">[</span><span class="n">o</span><span class="o">.</span><span class="n">_jOffsetRange</span><span class="p">(</span><span class="n">helper</span><span class="p">)</span> <span class="k">for</span> <span class="n">o</span> <span class="ow">in</span> <span class="n">offsetRanges</span><span class="p">]</span>
<span class="n">jleaders</span> <span class="o">=</span> <span class="nb">dict</span><span class="p">([(</span><span class="n">k</span><span class="o">.</span><span class="n">_jTopicAndPartition</span><span class="p">(</span><span class="n">helper</span><span class="p">),</span>
<span class="n">v</span><span class="o">.</span><span class="n">_jBroker</span><span class="p">(</span><span class="n">helper</span><span class="p">))</span> <span class="k">for</span> <span class="p">(</span><span class="n">k</span><span class="p">,</span> <span class="n">v</span><span class="p">)</span> <span class="ow">in</span> <span class="n">leaders</span><span class="o">.</span><span class="n">items</span><span class="p">()])</span>
<span class="n">jrdd</span> <span class="o">=</span> <span class="n">helper</span><span class="o">.</span><span class="n">createRDD</span><span class="p">(</span><span class="n">sc</span><span class="o">.</span><span class="n">_jsc</span><span class="p">,</span> <span class="n">kafkaParams</span><span class="p">,</span> <span class="n">joffsetRanges</span><span class="p">,</span> <span class="n">jleaders</span><span class="p">)</span>
<span class="k">except</span> <span class="n">Py4JJavaError</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
<span class="k">if</span> <span class="s">'ClassNotFoundException'</span> <span class="ow">in</span> <span class="nb">str</span><span class="p">(</span><span class="n">e</span><span class="o">.</span><span class="n">java_exception</span><span class="p">):</span>
<span class="n">KafkaUtils</span><span class="o">.</span><span class="n">_printErrorMsg</span><span class="p">(</span><span class="n">sc</span><span class="p">)</span>
<span class="k">raise</span> <span class="n">e</span>
<span class="n">ser</span> <span class="o">=</span> <span class="n">PairDeserializer</span><span class="p">(</span><span class="n">NoOpSerializer</span><span class="p">(),</span> <span class="n">NoOpSerializer</span><span class="p">())</span>
<span class="n">rdd</span> <span class="o">=</span> <span class="n">RDD</span><span class="p">(</span><span class="n">jrdd</span><span class="p">,</span> <span class="n">sc</span><span class="p">,</span> <span class="n">ser</span><span class="p">)</span><span class="o">.</span><span class="n">map</span><span class="p">(</span><span class="k">lambda</span> <span class="n">k_v</span><span class="p">:</span> <span class="p">(</span><span class="n">keyDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">0</span><span class="p">]),</span> <span class="n">valueDecoder</span><span class="p">(</span><span class="n">k_v</span><span class="p">[</span><span class="mi">1</span><span class="p">])))</span>
<span class="k">return</span> <span class="n">KafkaRDD</span><span class="p">(</span><span class="n">rdd</span><span class="o">.</span><span class="n">_jrdd</span><span class="p">,</span> <span class="n">rdd</span><span class="o">.</span><span class="n">ctx</span><span class="p">,</span> <span class="n">rdd</span><span class="o">.</span><span class="n">_jrdd_deserializer</span><span class="p">)</span>
</div>
<span class="nd">@staticmethod</span>
<span class="k">def</span> <span class="nf">_printErrorMsg</span><span class="p">(</span><span class="n">sc</span><span class="p">):</span>
<span class="k">print</span><span class="p">(</span><span class="s">"""</span>
<span class="s">________________________________________________________________________________________________</span>
<span class="s"> Spark Streaming's Kafka libraries not found in class path. Try one of the following.</span>
<span class="s"> 1. Include the Kafka library and its dependencies with in the</span>
<span class="s"> spark-submit command as</span>
<span class="s"> $ bin/spark-submit --packages org.apache.spark:spark-streaming-kafka:</span><span class="si">%s</span><span class="s"> ...</span>
<span class="s"> 2. Download the JAR of the artifact from Maven Central http://search.maven.org/,</span>
<span class="s"> Group Id = org.apache.spark, Artifact Id = spark-streaming-kafka-assembly, Version = </span><span class="si">%s</span><span class="s">.</span>
<span class="s"> Then, include the jar in the spark-submit command as</span>
<span class="s"> $ bin/spark-submit --jars <spark-streaming-kafka-assembly.jar> ...</span>
<span class="s">________________________________________________________________________________________________</span>
<span class="s">"""</span> <span class="o">%</span> <span class="p">(</span><span class="n">sc</span><span class="o">.</span><span class="n">version</span><span class="p">,</span> <span class="n">sc</span><span class="o">.</span><span class="n">version</span><span class="p">))</span>
</div>
<div class="viewcode-block" id="OffsetRange"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.OffsetRange">[docs]</a><span class="k">class</span> <span class="nc">OffsetRange</span><span class="p">(</span><span class="nb">object</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Represents a range of offsets from a single Kafka TopicAndPartition.</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">topic</span><span class="p">,</span> <span class="n">partition</span><span class="p">,</span> <span class="n">fromOffset</span><span class="p">,</span> <span class="n">untilOffset</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Create a OffsetRange to represent range of offsets</span>
<span class="sd"> :param topic: Kafka topic name.</span>
<span class="sd"> :param partition: Kafka partition id.</span>
<span class="sd"> :param fromOffset: Inclusive starting offset.</span>
<span class="sd"> :param untilOffset: Exclusive ending offset.</span>
<span class="sd"> """</span>
<span class="bp">self</span><span class="o">.</span><span class="n">topic</span> <span class="o">=</span> <span class="n">topic</span>
<span class="bp">self</span><span class="o">.</span><span class="n">partition</span> <span class="o">=</span> <span class="n">partition</span>
<span class="bp">self</span><span class="o">.</span><span class="n">fromOffset</span> <span class="o">=</span> <span class="n">fromOffset</span>
<span class="bp">self</span><span class="o">.</span><span class="n">untilOffset</span> <span class="o">=</span> <span class="n">untilOffset</span>
<span class="k">def</span> <span class="nf">__eq__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">other</span><span class="p">):</span>
<span class="k">if</span> <span class="nb">isinstance</span><span class="p">(</span><span class="n">other</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">__class__</span><span class="p">):</span>
<span class="k">return</span> <span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">topic</span> <span class="o">==</span> <span class="n">other</span><span class="o">.</span><span class="n">topic</span>
<span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">partition</span> <span class="o">==</span> <span class="n">other</span><span class="o">.</span><span class="n">partition</span>
<span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">fromOffset</span> <span class="o">==</span> <span class="n">other</span><span class="o">.</span><span class="n">fromOffset</span>
<span class="ow">and</span> <span class="bp">self</span><span class="o">.</span><span class="n">untilOffset</span> <span class="o">==</span> <span class="n">other</span><span class="o">.</span><span class="n">untilOffset</span><span class="p">)</span>
<span class="k">else</span><span class="p">:</span>
<span class="k">return</span> <span class="bp">False</span>
<span class="k">def</span> <span class="nf">__ne__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">other</span><span class="p">):</span>
<span class="k">return</span> <span class="ow">not</span> <span class="bp">self</span><span class="o">.</span><span class="n">__eq__</span><span class="p">(</span><span class="n">other</span><span class="p">)</span>
<span class="k">def</span> <span class="nf">__str__</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
<span class="k">return</span> <span class="s">"OffsetRange(topic: </span><span class="si">%s</span><span class="s">, partition: </span><span class="si">%d</span><span class="s">, range: [</span><span class="si">%d</span><span class="s"> -> </span><span class="si">%d</span><span class="s">]"</span> \
<span class="o">%</span> <span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">topic</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">partition</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">fromOffset</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">untilOffset</span><span class="p">)</span>
<span class="k">def</span> <span class="nf">_jOffsetRange</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">helper</span><span class="p">):</span>
<span class="k">return</span> <span class="n">helper</span><span class="o">.</span><span class="n">createOffsetRange</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">topic</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">partition</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">fromOffset</span><span class="p">,</span>
<span class="bp">self</span><span class="o">.</span><span class="n">untilOffset</span><span class="p">)</span>
</div>
<div class="viewcode-block" id="TopicAndPartition"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.TopicAndPartition">[docs]</a><span class="k">class</span> <span class="nc">TopicAndPartition</span><span class="p">(</span><span class="nb">object</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Represents a specific top and partition for Kafka.</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">topic</span><span class="p">,</span> <span class="n">partition</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Create a Python TopicAndPartition to map to the Java related object</span>
<span class="sd"> :param topic: Kafka topic name.</span>
<span class="sd"> :param partition: Kafka partition id.</span>
<span class="sd"> """</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_topic</span> <span class="o">=</span> <span class="n">topic</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_partition</span> <span class="o">=</span> <span class="n">partition</span>
<span class="k">def</span> <span class="nf">_jTopicAndPartition</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">helper</span><span class="p">):</span>
<span class="k">return</span> <span class="n">helper</span><span class="o">.</span><span class="n">createTopicAndPartition</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_topic</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">_partition</span><span class="p">)</span>
</div>
<div class="viewcode-block" id="Broker"><a class="viewcode-back" href="../../../pyspark.streaming.html#pyspark.streaming.kafka.Broker">[docs]</a><span class="k">class</span> <span class="nc">Broker</span><span class="p">(</span><span class="nb">object</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Represent the host and port info for a Kafka broker.</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">host</span><span class="p">,</span> <span class="n">port</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Create a Python Broker to map to the Java related object.</span>
<span class="sd"> :param host: Broker's hostname.</span>
<span class="sd"> :param port: Broker's port.</span>
<span class="sd"> """</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_host</span> <span class="o">=</span> <span class="n">host</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_port</span> <span class="o">=</span> <span class="n">port</span>
<span class="k">def</span> <span class="nf">_jBroker</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">helper</span><span class="p">):</span>
<span class="k">return</span> <span class="n">helper</span><span class="o">.</span><span class="n">createBroker</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_host</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">_port</span><span class="p">)</span>
</div>
<span class="k">class</span> <span class="nc">KafkaRDD</span><span class="p">(</span><span class="n">RDD</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> A Python wrapper of KafkaRDD, to provide additional information on normal RDD.</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">jrdd_deserializer</span><span class="p">):</span>
<span class="n">RDD</span><span class="o">.</span><span class="n">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">jrdd_deserializer</span><span class="p">)</span>
<span class="k">def</span> <span class="nf">offsetRanges</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Get the OffsetRange of specific KafkaRDD.</span>
<span class="sd"> :return: A list of OffsetRange</span>
<span class="sd"> """</span>
<span class="k">try</span><span class="p">:</span>
<span class="n">helperClass</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">ctx</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">java</span><span class="o">.</span><span class="n">lang</span><span class="o">.</span><span class="n">Thread</span><span class="o">.</span><span class="n">currentThread</span><span class="p">()</span><span class="o">.</span><span class="n">getContextClassLoader</span><span class="p">()</span> \
<span class="o">.</span><span class="n">loadClass</span><span class="p">(</span><span class="s">"org.apache.spark.streaming.kafka.KafkaUtilsPythonHelper"</span><span class="p">)</span>
<span class="n">helper</span> <span class="o">=</span> <span class="n">helperClass</span><span class="o">.</span><span class="n">newInstance</span><span class="p">()</span>
<span class="n">joffsetRanges</span> <span class="o">=</span> <span class="n">helper</span><span class="o">.</span><span class="n">offsetRangesOfKafkaRDD</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_jrdd</span><span class="o">.</span><span class="n">rdd</span><span class="p">())</span>
<span class="k">except</span> <span class="n">Py4JJavaError</span> <span class="k">as</span> <span class="n">e</span><span class="p">:</span>
<span class="k">if</span> <span class="s">'ClassNotFoundException'</span> <span class="ow">in</span> <span class="nb">str</span><span class="p">(</span><span class="n">e</span><span class="o">.</span><span class="n">java_exception</span><span class="p">):</span>
<span class="n">KafkaUtils</span><span class="o">.</span><span class="n">_printErrorMsg</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">ctx</span><span class="p">)</span>
<span class="k">raise</span> <span class="n">e</span>
<span class="n">ranges</span> <span class="o">=</span> <span class="p">[</span><span class="n">OffsetRange</span><span class="p">(</span><span class="n">o</span><span class="o">.</span><span class="n">topic</span><span class="p">(),</span> <span class="n">o</span><span class="o">.</span><span class="n">partition</span><span class="p">(),</span> <span class="n">o</span><span class="o">.</span><span class="n">fromOffset</span><span class="p">(),</span> <span class="n">o</span><span class="o">.</span><span class="n">untilOffset</span><span class="p">())</span>
<span class="k">for</span> <span class="n">o</span> <span class="ow">in</span> <span class="n">joffsetRanges</span><span class="p">]</span>
<span class="k">return</span> <span class="n">ranges</span>
<span class="k">class</span> <span class="nc">KafkaDStream</span><span class="p">(</span><span class="n">DStream</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> A Python wrapper of KafkaDStream</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">jdstream</span><span class="p">,</span> <span class="n">ssc</span><span class="p">,</span> <span class="n">jrdd_deserializer</span><span class="p">):</span>
<span class="n">DStream</span><span class="o">.</span><span class="n">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">jdstream</span><span class="p">,</span> <span class="n">ssc</span><span class="p">,</span> <span class="n">jrdd_deserializer</span><span class="p">)</span>
<span class="k">def</span> <span class="nf">foreachRDD</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">func</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Apply a function to each RDD in this DStream.</span>
<span class="sd"> """</span>
<span class="k">if</span> <span class="n">func</span><span class="o">.</span><span class="n">__code__</span><span class="o">.</span><span class="n">co_argcount</span> <span class="o">==</span> <span class="mi">1</span><span class="p">:</span>
<span class="n">old_func</span> <span class="o">=</span> <span class="n">func</span>
<span class="n">func</span> <span class="o">=</span> <span class="k">lambda</span> <span class="n">r</span><span class="p">,</span> <span class="n">rdd</span><span class="p">:</span> <span class="n">old_func</span><span class="p">(</span><span class="n">rdd</span><span class="p">)</span>
<span class="n">jfunc</span> <span class="o">=</span> <span class="n">TransformFunction</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_sc</span><span class="p">,</span> <span class="n">func</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">_jrdd_deserializer</span><span class="p">)</span> \
<span class="o">.</span><span class="n">rdd_wrapper</span><span class="p">(</span><span class="k">lambda</span> <span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">ser</span><span class="p">:</span> <span class="n">KafkaRDD</span><span class="p">(</span><span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">ser</span><span class="p">))</span>
<span class="n">api</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">_ssc</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">PythonDStream</span>
<span class="n">api</span><span class="o">.</span><span class="n">callForeachRDD</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_jdstream</span><span class="p">,</span> <span class="n">jfunc</span><span class="p">)</span>
<span class="k">def</span> <span class="nf">transform</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">func</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Return a new DStream in which each RDD is generated by applying a function</span>
<span class="sd"> on each RDD of this DStream.</span>
<span class="sd"> `func` can have one argument of `rdd`, or have two arguments of</span>
<span class="sd"> (`time`, `rdd`)</span>
<span class="sd"> """</span>
<span class="k">if</span> <span class="n">func</span><span class="o">.</span><span class="n">__code__</span><span class="o">.</span><span class="n">co_argcount</span> <span class="o">==</span> <span class="mi">1</span><span class="p">:</span>
<span class="n">oldfunc</span> <span class="o">=</span> <span class="n">func</span>
<span class="n">func</span> <span class="o">=</span> <span class="k">lambda</span> <span class="n">t</span><span class="p">,</span> <span class="n">rdd</span><span class="p">:</span> <span class="n">oldfunc</span><span class="p">(</span><span class="n">rdd</span><span class="p">)</span>
<span class="k">assert</span> <span class="n">func</span><span class="o">.</span><span class="n">__code__</span><span class="o">.</span><span class="n">co_argcount</span> <span class="o">==</span> <span class="mi">2</span><span class="p">,</span> <span class="s">"func should take one or two arguments"</span>
<span class="k">return</span> <span class="n">KafkaTransformedDStream</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">func</span><span class="p">)</span>
<span class="k">class</span> <span class="nc">KafkaTransformedDStream</span><span class="p">(</span><span class="n">TransformedDStream</span><span class="p">):</span>
<span class="sd">"""</span>
<span class="sd"> Kafka specific wrapper of TransformedDStream to transform on Kafka RDD.</span>
<span class="sd"> """</span>
<span class="k">def</span> <span class="nf">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">prev</span><span class="p">,</span> <span class="n">func</span><span class="p">):</span>
<span class="n">TransformedDStream</span><span class="o">.</span><span class="n">__init__</span><span class="p">(</span><span class="bp">self</span><span class="p">,</span> <span class="n">prev</span><span class="p">,</span> <span class="n">func</span><span class="p">)</span>
<span class="nd">@property</span>
<span class="k">def</span> <span class="nf">_jdstream</span><span class="p">(</span><span class="bp">self</span><span class="p">):</span>
<span class="k">if</span> <span class="bp">self</span><span class="o">.</span><span class="n">_jdstream_val</span> <span class="ow">is</span> <span class="ow">not</span> <span class="bp">None</span><span class="p">:</span>
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">_jdstream_val</span>
<span class="n">jfunc</span> <span class="o">=</span> <span class="n">TransformFunction</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">_sc</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">func</span><span class="p">,</span> <span class="bp">self</span><span class="o">.</span><span class="n">prev</span><span class="o">.</span><span class="n">_jrdd_deserializer</span><span class="p">)</span> \
<span class="o">.</span><span class="n">rdd_wrapper</span><span class="p">(</span><span class="k">lambda</span> <span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">ser</span><span class="p">:</span> <span class="n">KafkaRDD</span><span class="p">(</span><span class="n">jrdd</span><span class="p">,</span> <span class="n">ctx</span><span class="p">,</span> <span class="n">ser</span><span class="p">))</span>
<span class="n">dstream</span> <span class="o">=</span> <span class="bp">self</span><span class="o">.</span><span class="n">_sc</span><span class="o">.</span><span class="n">_jvm</span><span class="o">.</span><span class="n">PythonTransformedDStream</span><span class="p">(</span><span class="bp">self</span><span class="o">.</span><span class="n">prev</span><span class="o">.</span><span class="n">_jdstream</span><span class="o">.</span><span class="n">dstream</span><span class="p">(),</span> <span class="n">jfunc</span><span class="p">)</span>
<span class="bp">self</span><span class="o">.</span><span class="n">_jdstream_val</span> <span class="o">=</span> <span class="n">dstream</span><span class="o">.</span><span class="n">asJavaDStream</span><span class="p">()</span>
<span class="k">return</span> <span class="bp">self</span><span class="o">.</span><span class="n">_jdstream_val</span>
</pre></div>
</div>
</div>
</div>
<div class="sphinxsidebar" role="navigation" aria-label="main navigation">
<div class="sphinxsidebarwrapper">
<p class="logo"><a href="../../../index.html">
<img class="logo" src="../../../_static/spark-logo-hd.png" alt="Logo"/>
</a></p>
<div id="searchbox" style="display: none" role="search">
<h3>Quick search</h3>
<form class="search" action="../../../search.html" method="get">
<input type="text" name="q" />
<input type="submit" value="Go" />
<input type="hidden" name="check_keywords" value="yes" />
<input type="hidden" name="area" value="default" />
</form>
<p class="searchtip" style="font-size: 90%">
Enter search terms or a module, class or function name.
</p>
</div>
<script type="text/javascript">$('#searchbox').show(0);</script>
</div>
</div>
<div class="clearer"></div>
</div>
<div class="related" role="navigation" aria-label="related navigation">
<h3>Navigation</h3>
<ul>
<li class="nav-item nav-item-0"><a href="../../../index.html">PySpark 1.5.0 documentation</a> »</li>
<li class="nav-item nav-item-1"><a href="../../index.html" >Module code</a> »</li>
</ul>
</div>
<div class="footer" role="contentinfo">
© Copyright .
Created using <a href="http://sphinx-doc.org/">Sphinx</a> 1.3.1.
</div>
</body>
</html>