PySpark多流式数据集左外连接时如何保留null值
我来帮你捋捋这个问题哈,你说第一次左外连接还能看到未匹配的null行,第二次连完就没了,主要有两个核心原因,咱们一个个说清楚:
1. 字段冲突+select语句的坑
你第一次用unigue_ids左连df_b后,结果里的some_parameter是来自df_b的;第二次左连df_a时,df_a也有个some_parameter(值是"A"),这时候结果里会有两个同名的some_parameter字段,一个来自B,一个来自A。但你最后select只写了some_parameter,PySpark只会挑其中一个(一般是最后join进来的那个),更关键的是:那些两次都没匹配到的行,本来应该保留,但可能因为你没正确处理字段,或者是下面要说的流式状态问题,导致被过滤掉了。
2. 流式Join的水印状态问题
PySpark流式Join靠水印管理状态,防止状态无限膨胀。你第一次用的unigue_ids和df_b都是带水印的,没问题;但第一次连完后的DataFrame没重新设置水印,第二次用这个无水印的DF去连df_a时,Spark没法正确判断哪些数据是有效状态,可能会把那些未匹配的null行当成过期数据清掉,所以你最后就看不到它们了。
另外哦,你的get_join_str函数里有个笔误:前面代码里用的是some_id,后面函数里写成networkcallref了,这也会导致连接条件不匹配,得先修正这个!
解决办法来了,你可以试试这几个方案:
方案一:每次连接后重新设置水印
每次左连完,都给新DF加上基于timestamp的水印,让Spark能正确管理状态:
# 第一次左连+重新加水印 df = unigue_ids.alias("left") \ .join(df_b, F.expr(get_join_str("B")), how="leftOuter") \ .withWatermark("timestamp", watermark) # 重新设置水印 # 第二次左连+加水印 df = df.alias("left") \ .join(df_a, F.expr(get_join_str("A")), how="leftOuter") \ .withWatermark("timestamp", watermark) # 最后select的时候,用coalesce合并两个some_parameter,优先取B的,没有就取A的,都没有就是null df = df.select( 'some_id', 'timestamp', F.coalesce(F.col("B.some_parameter"), F.col("A.some_parameter")).alias("some_parameter") )
方案二:一次性完成所有左连(更推荐)
流式处理里多次Join会增加状态复杂度,不如把所有要连的DF一次性连完,Spark处理起来更高效:
# 一次性连df_b和df_a df = unigue_ids.alias("left") \ .join(df_b, F.expr(get_join_str("B")), how="leftOuter") \ .join(df_a, F.expr(get_join_str("A")), how="leftOuter") \ .select( 'some_id', 'timestamp', F.coalesce(F.col("B.some_parameter"), F.col("A.some_parameter")).alias("some_parameter") )
这样既能避免状态问题,也能保留所有未匹配的行。
方案三:修正连接条件的字段名
把get_join_str里的networkcallref改成some_id,确保连接条件和你的字段一致:
def get_join_str(some_parameter): return f"""left.some_id = {some_parameter}.some_id AND {some_parameter}.timestamp BETWEEN left.timestamp AND left.timestamp + interval 1 second """
按上面的方法改完,那些两次都没匹配到的行应该就能保留下来,some_parameter会显示null啦~
备注:内容来源于stack exchange,提问作者GFR

