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

如何将Schema不同的多个CSV文件导入单个Spark DataFrame?

解决Spark加载Schema差异大的CSV文件并统计总词数

问题原因

Spark默认读取多个CSV时,若未开启mergeSchema参数,会采用推断出的单一Schema(通常基于采样或某一个文件的结构),导致Schema不匹配的文件内容无法正确加载,出现仅显示单个文件Schema的情况。

解决步骤

1. 验证两个文件的结构

先分别读取两个CSV,确认各自的列结构及文本内容所在列:

# 读取reddit数据集,查看结构与样本内容
df_reddit = spark.read.format("csv").option("header", "false").load("abfss://lmne.dfs.core.windows.net/csvs/MachineLearning_reddit.csv")
df_reddit.printSchema()
df_reddit.show(5)

# 读取bbc新闻数据集,查看结构与样本内容
df_bbc = spark.read.format("csv").option("header", "false").load("abfss://test1@lmne.dfs.core.windows.net/csvs/bbc_news.csv")
df_bbc.printSchema()
df_bbc.show(5)

2. 统一Schema后合并数据集

根据上一步的结果,提取两个数据集中的文本列,重命名为统一列名后合并:

# 示例:假设reddit的文本在_c0列,bbc的文本在_c1列,统一重命名为text
df_reddit_text = df_reddit.select("_c0").withColumnRenamed("_c0", "text")
df_bbc_text = df_bbc.select("_c1").withColumnRenamed("_c1", "text")

# 合并两个数据集
combined_df = df_reddit_text.unionAll(df_bbc_text)

如果需要保留所有列并自动合并Schema,可开启mergeSchema参数:

df_combined = spark.read.format("csv")\
    .option("header", "false")\
    .option("mergeSchema", "true")\
    .load(paths)

此方式会保留所有列,缺失的列填充null,后续可通过concat_ws合并所有文本列。

3. 统计总词数

使用Spark内置函数拆分文本、展开单词后统计总数:

from pyspark.sql.functions import split, explode, col, concat_ws

# 针对统一text列的数据集
word_df = combined_df.select(split(col("text"), "\\s+").alias("words"))\
    .select(explode(col("words")).alias("word"))

# 针对mergeSchema后的多列数据集,先合并所有列文本
# word_df = df_combined.select(concat_ws(" ", *df_combined.columns).alias("text"))\
#     .select(split(col("text"), "\\s+").alias("words"))\
#     .select(explode(col("words")).alias("word"))

# 统计总词数
total_words = word_df.count()
print(f"总词数:{total_words}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 17:35:26