Spark Structured Streaming多流左外连接失效问题求助
Spark Structured Streaming三流左外连接异常排查求助
问题现象
- 两DataFrame流左外连接可在水印过期后正常输出null值,但三流连接时出现异常
- 三个流均从Kafka读取(session、comparison、filter),预期水印过期后,若comparison或filter流无数据,应输出session数据搭配null值,但session数据丢失,即使等待10分钟也未在控制台出现
- 未找到三流流连接的相关资料,求助排查问题或确认Spark Structured Streaming是否存在三流及以上流连接的限制
相关代码
watermark_duration = "5 minutes" # 三个流共用该水印时长 interval_value = "5 minutes" # Session流处理 session_df = spark.readStream.format("kafka").options(**options).option("subscribe", session_source_topic).load() session_df = session_df.selectExpr("CAST(value AS STRING) as session_raw_data", "Timestamp as session_Timestamp")\ .withColumn("session_data", split("session_raw_data", "\\|"))\ .selectExpr("session_data", "session_Timestamp") session_df = session_df.select( col("session_data")[0].alias("id"), col("session_data")[1].alias("session_id"), col("session_Timestamp") ).withWatermark("session_Timestamp", watermark_duration) # Comparison流处理 comparison_df = spark.readStream.format("kafka").options(**options).option("subscribe", comparison_source_topic).load() comparison_df = comparison_df.selectExpr("CAST(value AS STRING) as comparison_raw_data", "Timestamp as comparison_Timestamp")\ .withColumn("comparison_data", split("comparison_raw_data", "\\|"))\ .selectExpr("comparison_data", "comparison_Timestamp") comparison_df = comparison_df.select( col("comparison_data")[0].alias("comparison_id"), col("comparison_data")[1].alias("comparison_session_id"), col("comparison_Timestamp") ).withWatermark("watermark_duration", comparison_watermark) # 此处存在参数顺序错误 # Filter流处理 filter_df = spark.readStream.format("kafka").options(**options).option("subscribe", filter_source_topic).load() filter_df = filter_df.selectExpr("CAST(value AS STRING) as filter_raw_data", "Timestamp as filter_Timestamp_Timestamp")\ .withColumn("filter_data", split("filter_raw_data", "\\|"))\ .selectExpr("filter_data", "filter_Timestamp_Timestamp") filter_df = filter_df.select( col("filter_data")[0].alias("filter_id"), col("filter_data")[1].alias("filter_session_id"), col("filter_Timestamp_Timestamp") ).withWatermark("filter_Timestamp_Timestamp", watermark_duration) # 连接条件 comparison_join_expr = "session_id=comparison_session_id AND " + \ " comparison_Timestamp between session_Timestamp and session_Timestamp + interval " + interval_value filter_join_expr = "session_id=filter_session_id AND " + \ " filter_Timestamp_Timestamp between session_Timestamp and session_Timestamp + interval " + interval_value # 三流左外连接 joined_df1 = session_df.join(comparison_df, expr(comparison_join_expr), "leftOuter")\ .join(filter_df, expr(filter_join_expr), "leftOuter").drop("filter_session_id") # 测试输出 session_df.writeStream.outputMode("append").format("console").start() comparison_df.writeStream.outputMode("append").format("console").start() filter_df.writeStream.outputMode("append").format("console").start() joined_df1.writeStream.outputMode("append").format("console").start() # 注:原代码中joined_df2、joined_df3未定义,此处忽略 spark.streams.awaitAnyTermination()
预期输出
当comparison_df和filter_df无数据时,输出session_df的完整数据,对应comparison和filter相关字段为null值,示例如下:
| id | session_id | session_Timestamp | comparison_id | comparison_session_id | comparison_Timestamp | filter_id | filter_Timestamp_Timestamp |
|---|---|---|---|---|---|---|---|
| 1 | s_001 | 2024-05-20 10:00:00 | null | null | null | null | null |
排查建议
- 修正水印参数错误:comparison_df的
withWatermark参数顺序写反,正确应为withWatermark("comparison_Timestamp", watermark_duration),错误的配置会导致comparison流的水印不生效,无法触发连接后的输出 - 修正时间列命名:filter_df的时间列命名为
filter_Timestamp_Timestamp存在重复,建议改为filter_Timestamp,避免后续逻辑中可能出现的解析问题 - 检查连接触发逻辑:多流左外连接依赖主流水印触发输出,需确保每个连接的事件时间条件(
between区间)和水印配置匹配,当后续流无数据时,主流水印过期后才会输出匹配null的结果 - 验证Spark版本:旧版本Spark(如3.0以下)在多流左外连接处理上可能存在bug,建议升级到3.2及以上稳定版本
- 确认输出模式:append模式下,数据只有在水印过期、确认无后续匹配数据时才会输出,需等待超过
watermark_duration + interval_value的时长再观察结果
内容的提问来源于stack exchange,提问作者Paulnaveen
相关产品推荐
相关产品推荐

