-
- Downloads
[SPARK-5155] [PYSPARK] [STREAMING] Mqtt streaming support in Python
This PR is based on #4229, thanks prabeesh. Closes #4229 Author: Prabeesh K <prabsmails@gmail.com> Author: zsxwing <zsxwing@gmail.com> Author: prabs <prabsmails@gmail.com> Author: Prabeesh K <prabeesh.k@namshi.com> Closes #7833 from zsxwing/pr4229 and squashes the following commits: 9570bec [zsxwing] Fix the variable name and check null in finally 4a9c79e [zsxwing] Fix pom.xml indentation abf5f18 [zsxwing] Merge branch 'master' into pr4229 935615c [zsxwing] Fix the flaky MQTT tests 47278c5 [zsxwing] Include the project class files 478f844 [zsxwing] Add unpack 5f8a1d4 [zsxwing] Make the maven build generate the test jar for Python MQTT tests 734db99 [zsxwing] Merge branch 'master' into pr4229 126608a [Prabeesh K] address the comments b90b709 [Prabeesh K] Merge pull request #1 from zsxwing/pr4229 d07f454 [zsxwing] Register StreamingListerner before starting StreamingContext; Revert unncessary changes; fix the python unit test a6747cb [Prabeesh K] wait for starting the receiver before publishing data 87fc677 [Prabeesh K] address the comments: 97244ec [zsxwing] Make sbt build the assembly test jar for streaming mqtt 80474d1 [Prabeesh K] fix 1f0cfe9 [Prabeesh K] python style fix e1ee016 [Prabeesh K] scala style fix a5a8f9f [Prabeesh K] added Python test 9767d82 [Prabeesh K] implemented Python-friendly class a11968b [Prabeesh K] fixed python style 795ec27 [Prabeesh K] address comments ee387ae [Prabeesh K] Fix assembly jar location of mqtt-assembly 3f4df12 [Prabeesh K] updated version b34c3c1 [prabs] adress comments 3aa7fff [prabs] Added Python streaming mqtt word count example b7d42ff [prabs] Mqtt streaming support in Python
Showing
- dev/run-tests.py 2 additions, 0 deletionsdev/run-tests.py
- dev/sparktestsupport/modules.py 2 additions, 0 deletionsdev/sparktestsupport/modules.py
- docs/streaming-programming-guide.md 1 addition, 1 deletiondocs/streaming-programming-guide.md
- examples/src/main/python/streaming/mqtt_wordcount.py 58 additions, 0 deletionsexamples/src/main/python/streaming/mqtt_wordcount.py
- external/mqtt-assembly/pom.xml 102 additions, 0 deletionsexternal/mqtt-assembly/pom.xml
- external/mqtt/pom.xml 28 additions, 0 deletionsexternal/mqtt/pom.xml
- external/mqtt/src/main/assembly/assembly.xml 44 additions, 0 deletionsexternal/mqtt/src/main/assembly/assembly.xml
- external/mqtt/src/main/scala/org/apache/spark/streaming/mqtt/MQTTUtils.scala 16 additions, 0 deletions...ain/scala/org/apache/spark/streaming/mqtt/MQTTUtils.scala
- external/mqtt/src/test/scala/org/apache/spark/streaming/mqtt/MQTTStreamSuite.scala 15 additions, 103 deletions...ala/org/apache/spark/streaming/mqtt/MQTTStreamSuite.scala
- external/mqtt/src/test/scala/org/apache/spark/streaming/mqtt/MQTTTestUtils.scala 111 additions, 0 deletions...scala/org/apache/spark/streaming/mqtt/MQTTTestUtils.scala
- pom.xml 1 addition, 0 deletionspom.xml
- project/SparkBuild.scala 9 additions, 3 deletionsproject/SparkBuild.scala
- python/pyspark/streaming/mqtt.py 72 additions, 0 deletionspython/pyspark/streaming/mqtt.py
- python/pyspark/streaming/tests.py 104 additions, 2 deletionspython/pyspark/streaming/tests.py
Loading
Please register or sign in to comment