Stream-Stream Left Outer Join in PySpark Structured Streaming causing Java OOM errors

Viewed 87

Consider the following simplified code:

left_stream= spark.readStream. ...
right_stream = spark.readStream. ...

# Apply watermarks on event-time columns
left_stream_with_watermark = left_stream.withWatermark("leftTime", "20 seconds")
right_stream_with_watermark = right_stream.withWatermark("rightTime", "20 seconds")

# Join with event-time constraints
left_stream_with_watermark.join(
right_stream_with_watermark,
expr("""
leftId = rightId AND
leftTime >= rightTime - interval 20 seconds AND
leftTime <= rightTime + interval 20 seconds
"""),
"leftOuter"
)

On the surface the ON clause appears valid for this join and the right side rows should expire from state if they are not within 60 seconds of the left side watermark.

In practice, we observe the right side state store for this streaming query increasing in size over time (as if right side rows are not expiring). And as a result Java OOM errors occur frequently on executors running the related tasks.

Expectation is that records on the right side of the join should expire, but they in fact remain in right side state store.

0 Answers
Related