-
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathstructured_streaming_example.py
More file actions
31 lines (26 loc) · 1 KB
/
Copy pathstructured_streaming_example.py
File metadata and controls
31 lines (26 loc) · 1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
# Minimal PySpark Structured Streaming example.
# Requires the connector JAR on the driver/executor classpath, e.g.:
# pyspark --packages io.github.juarezr:spark-streaming-google-pubsub_2.12:0.9.4
# Spark 4.x (Scala 2.13):
# pyspark --packages io.github.juarezr:spark-streaming-google-pubsub_2.13:0.9.4
# or:
# spark-submit --packages ... examples/python/structured_streaming_example.py
from pyspark.sql import SparkSession
from pyspark.sql.streaming import Trigger
spark = SparkSession.builder.appName("pubsub-pyspark-example").getOrCreate()
messages = (
spark.readStream.format("google-pubsub")
.option("projectId", "my-project")
.option("subscription", "my-subscription")
.option("ackMode", "afterCommit")
.option("gatherMode", "batch")
.load()
)
query = (
messages.writeStream.format("console")
.option("truncate", "false")
.option("checkpointLocation", "/tmp/pubsub-pyspark-checkpoint")
.trigger(Trigger.ProcessingTime("1 second"))
.start()
)
query.awaitTermination()