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

如何在Google Colab中显示PySpark流计算的聚合结果?

问题排查与解决方案

1. 确认Drive挂载与路径正确性

  • 先执行!ls "/content/drive/My Drive/Pyspark/"验证目标文件夹是否存在,以及是否有CSV文件(或后续会添加文件)。若路径不匹配,直接修正load()中的路径即可。
  • 注意:代码中路径的空格已用引号包裹,这部分没问题,但要保证文件夹名称和路径完全一致。

2. 适配PySpark流对Drive的文件监控

由于Colab挂载的Drive是FUSE文件系统,PySpark默认的文件系统通知机制无法实时感知文件变化,需要添加参数主动扫描文件:

  • 加入option("maxFilesPerTrigger", 1),限制每次触发仅处理一个新文件;
  • 加入option("cleanSource", "archive")搭配archivePath,将已处理文件归档到指定文件夹,避免重复处理;
  • 可选添加option("fileCreationTime", "latest"),确保只处理新创建的文件。

修改后的读取代码片段:

df = spark.readStream.format("csv")\
    .schema(schema)\
    .option("header", True)\
    .option("sep", ",")\
    .option("maxFilesPerTrigger", 1)\
    .option("cleanSource", "archive")\
    .option("archivePath", "/content/drive/My Drive/Pyspark/processed/")\
    .load("/content/drive/My Drive/Pyspark/")

提前创建processed文件夹用于存放已处理文件

3. 优化聚合结果与输出模式

  • 原聚合代码生成的列名为sum(sales),可以添加别名提升可读性:
from pyspark.sql.functions import sum
df = df.groupBy("Shop").agg(sum("Sales").alias("Total_Sales"))
  • 若使用update模式,只有聚合结果变化时才会输出;如果希望看到所有历史聚合结果,改用complete模式(你的groupby聚合符合complete模式的要求),同时添加truncate=False避免内容被截断:
df.writeStream.format("console")\
    .outputMode("complete")\
    .option("truncate", False)\
    .start().awaitTermination()

4. 查看Spark运行日志

Colab默认Spark日志级别为WARN,无法看到流处理的触发细节,调整为INFO级别可以确认流查询是否正常启动:

spark.sparkContext.setLogLevel("INFO")

运行后会看到类似Streaming query started的日志,验证流任务处于运行状态。

5. 正确添加测试文件

  • 不要直接在Colab文件管理器上传文件到Drive路径,建议通过Drive网页端上传,或用Colab的文件上传工具直接传到挂载的文件夹;
  • 文件上传完成后等待1-2分钟,Drive挂载存在同步延迟,PySpark需要时间扫描到新文件。

完整修正后的代码

from google.colab import drive
drive.mount('/content/drive')

# 提前创建归档文件夹
import os
processed_path = "/content/drive/My Drive/Pyspark/processed/"
if not os.path.exists(processed_path):
    os.makedirs(processed_path)

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType,StructField,IntegerType,StringType
from pyspark.sql.functions import sum

spark = SparkSession.builder.master("local[*]").getOrCreate()
spark.sparkContext.setLogLevel("INFO")

schema=StructType(
    [
        StructField('File',StringType(),True),
        StructField('Shop',StringType(),True),
        StructField('Sales',IntegerType(),True)
    ]
)

df = spark.readStream.format("csv")\
    .schema(schema)\
    .option("header", True)\
    .option("sep", ",")\
    .option("maxFilesPerTrigger", 1)\
    .option("cleanSource", "archive")\
    .option("archivePath", processed_path)\
    .load("/content/drive/My Drive/Pyspark/")

# 聚合并设置别名
df = df.groupBy("Shop").agg(sum("Sales").alias("Total_Sales"))

# 输出到控制台
df.writeStream.format("console")\
    .outputMode("complete")\
    .option("truncate", False)\
    .start().awaitTermination()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:45:23