Created
August 5, 2021 19:23
-
-
Save ychennay/afc79f0a03a25cee562b1fa114303cf0 to your computer and use it in GitHub Desktop.
Spark Structured Streaming Kinesis Watermarks
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
from pyspark.sql.functions import window | |
# configure reading from the stream | |
kinesis_df = spark.readStream.format("kinesis") | |
.option("streamName", KINESIS_STREAM_NAME) | |
.option("region", AWS_REGION) | |
.option("roleArn", KINESIS_ACCESS_ROLE_ARN | |
.option("initialPosition", "latest") | |
.load() | |
# configure custom logic for watermarks / event time windows | |
kinesis_df.withWatermark("eventTime", "10 minutes") \ | |
.groupBy( | |
"userId", | |
window("eventTime", "10 minutes", "5 minutes")) |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment