You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark Streaming双流关联后实现单词累计求和的问题

PySpark Streaming跨窗口累计关联聚合实现

我正在使用PySpark Streaming对两个数据流进行关联与聚合:一个数据流包含word列,另一个数据流包含word对应的整数num,需要输出每个word对应的num累计总和。

测试代码如下(实际场景使用Kafka,示例采用socket流):

# 可自定义流配置,实际场景使用Kafka
# 示例用终端nc命令输入数据:nc -lk 9999
s1 = (
    spark.readStream.format("socket")
                    .options(host="localhost", port=9999)
                    .load()
)

# nc -lk 8888
s2 = (
    spark.readStream.format("socket")
                    .options(host="localhost", port=8888)
                    .load()
)

# 命名列并添加水印所需的时间窗口
w1 = s1.select(
    fn.col("value").alias("word"),
    fn.window(fn.expr("current_timestamp"), "30 seconds").alias("ts")
).withWatermark("ts", "10 seconds")

w2 = s2.select(
    fn.col("value").alias("word"),
    fn.window(fn.expr("current_timestamp"), "30 seconds").alias("ts"),
    fn.expr("int(rand() * 10)").alias("num")
).withWatermark("ts", "10 seconds")

# 关联并聚合
j = w1.join(
    w2, 
    ["ts", "word"], 
    "left"
).groupBy("word").sum("num")

# 报错:关联操作不支持"complete"输出模式
q = j.writeStream.outputMode("complete").foreachBatch(lambda df, _: df.show(truncate=False)).start()

期望效果

处理第一个时间窗口数据后,输出各word的num求和结果;处理第二个窗口数据后,输出跨窗口的累计总和:

第一个窗口数据

时间wordnum
t1foo1
t1foo2
t1bar3
t1bar4

第一个窗口输出

wordsum(num)
foo3
bar7

第二个窗口数据

时间wordnum
t2foo3
t2foo4
t2bar5
t2bar5

第二个窗口输出(跨窗口累计)

wordsum(num)
foo10
bar17

报错原因

带水印的流join操作不支持complete输出模式——join后的流基于窗口增量数据,complete模式要求维护全量聚合状态,但带水印的流会自动清理过期状态,两者逻辑冲突。

解决方案

要实现跨窗口累计总和,可通过以下两种方式调整:

方案一:直接按word关联+累计状态维护(推荐,匹配需求)

放弃窗口对齐逻辑,仅按word关联数据流,通过foreachBatch维护全局累计状态:

from pyspark.sql import functions as fn

# 处理s1:提取目标word(若需过滤可在此添加逻辑)
w1 = s1.select(fn.col("value").alias("word"))

# 处理s2:提取word和对应num值
w2 = s2.select(
    fn.col("value").alias("word"),
    fn.expr("int(rand() * 10)").alias("num")
)

# 按word关联两个流
j = w1.join(w2, on="word", how="left")

# 初始化累计状态(生产环境建议用Redis/HBase等外部存储替代全局变量)
total_sums = {}

def update_total(df, batch_id):
    global total_sums
    # 计算当前批次的word求和结果
    batch_sum = df.groupBy("word").sum("num").na.fill(0).collect()
    # 更新累计总和
    for row in batch_sum:
        word = row["word"]
        batch_val = row["sum(num)"]
        total_sums[word] = total_sums.get(word, 0) + batch_val
    # 转换为DataFrame展示结果
    result_df = spark.createDataFrame(
        [(k, v) for k, v in total_sums.items()],
        ["word", "sum(num)"]
    )
    result_df.show(truncate=False)

# 使用update输出模式处理增量数据
q = j.writeStream.outputMode("update").foreachBatch(update_total).start()

q.awaitTermination()

方案二:保留窗口逻辑+累计聚合

若必须保留窗口对齐关联,先对窗口内数据聚合,再合并到累计状态:

from pyspark.sql import functions as fn

# 处理s1:添加窗口和水印
w1 = s1.select(
    fn.col("value").alias("word"),
    fn.window(fn.expr("current_timestamp"), "30 seconds").alias("ts")
).withWatermark("ts", "10 seconds")

# 处理s2:添加窗口、水印,先聚合窗口内的num总和
w2 = s2.select(
    fn.col("value").alias("word"),
    fn.window(fn.expr("current_timestamp"), "30 seconds").alias("ts"),
    fn.expr("int(rand() * 10)").alias("num")
).withWatermark("ts", "10 seconds") \
 .groupBy("ts", "word").sum("num")

# 关联同一窗口内的word数据
j = w1.join(w2, ["ts", "word"], "left").select("word", "sum(num)").na.fill(0)

# 初始化累计状态
total_sums = {}

def update_total(df, batch_id):
    global total_sums
    # 计算当前批次窗口聚合后的结果
    batch_sum = df.groupBy("word").sum("sum(num)").collect()
    # 更新累计总和
    for row in batch_sum:
        word = row["word"]
        batch_val = row["sum(sum(num))"]
        total_sums[word] = total_sums.get(word, 0) + batch_val
    # 转换为DataFrame展示
    result_df = spark.createDataFrame(
        [(k, v) for k, v in total_sums.items()],
        ["word", "sum(num)"]
    )
    result_df.show(truncate=False)

# 使用append输出模式处理窗口增量数据
q = j.writeStream.outputMode("append").foreachBatch(update_total).start()

q.awaitTermination()

关键说明

  • 生产环境建议用外部存储(如Redis、HBase)维护累计状态,避免全局变量的分布式一致性问题;
  • update/append是流join支持的输出模式,complete仅适用于无水印的全量聚合场景;
  • 若s1用于过滤目标word,可先对s1做去重处理,避免重复关联导致重复计算。

内容的提问来源于stack exchange,提问作者RSHAP

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 09:27:06