Spark Structured Streaming 2.3.0双流左外连接未输出Null值问题求助
Ah, I remember this common gotcha with Spark 2.3.0's structured streaming left outer joins—you're running into the issue where the engine holds onto left stream records indefinitely, waiting for matching right stream records, instead of outputting the null-filled rows for unmatched entries. Let's walk through exactly what you need to adjust to get your desired behavior.
Key Issues in Your Current Setup
From your snippet, I notice two critical missing pieces that prevent the null records from being emitted:
- No watermark on the right stream (
ds2): Both streams need watermarks to help Spark determine when it's safe to conclude no future matching records will arrive. - Missing time-range constraint in the join condition: Spark's streaming left outer join requires a time-based window to bound the potential matches from the right stream. Without this, the engine can't know when to stop waiting for matches and output the null rows.
Step-by-Step Solution
1. Add Watermark to Both Streams
First, ensure both ds1 and ds2 have watermarks defined. Your ds1 already has one, but you need to add the same (or appropriate) watermark to ds2:
val ds1Map = ds1 .selectExpr("Id AS ds1_Id", "ds1_timestamp") .withWatermark("ds1_timestamp", "10 seconds") // Add watermark to ds2Map val ds2Map = ds2 .selectExpr("Id AS ds2_Id", "ds2_timestamp") .withWatermark("ds2_timestamp", "10 seconds") // Critical addition!
2. Include a Time-Range Condition in the Join
Next, your join condition must include a time-based constraint to tell Spark how far back/forward to look for matches in ds2. For example, if you want to match ds1 records with ds2 records that fall within a 10-second window around ds1_timestamp, use this condition:
import org.apache.spark.sql.functions.expr val joinedStream = ds1Map.join( ds2Map, // Join on matching IDs + time window constraint expr("ds1_Id = ds2_Id AND ds1_timestamp BETWEEN ds2_timestamp - interval 10 seconds AND ds2_timestamp + interval 10 seconds"), joinType = "leftOuter" )
3. Use the Correct Output Mode
Finally, use the append output mode (this is the mode that will emit records once their state is finalized, i.e., when Spark is sure no more matches will arrive):
joinedStream.writeStream .outputMode("append") .format("console") .start() .awaitTermination()
Why This Works
- Watermarks: The watermark on
ds1tells Spark that once the current processing time passesds1_timestamp + 10 seconds, no olderds1records will arrive. Theds2watermark does the same for the right stream. Together, these let Spark safely clean up state and emit unmatchedds1records (with nulls fords2fields) once the window for matching has passed. - Time-Range Constraint: This bounds the set of potential
ds2records that can match a givends1record. Without it, Spark has no way to know when to stop waiting for matches, so it holds ontods1records forever.
Keep in mind that Spark 2.3.0 has stricter requirements for streaming joins compared to newer versions—this setup is mandatory to get left outer joins to emit null records as expected.
内容的提问来源于stack exchange,提问作者user3191291

