如何将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
相关产品推荐
相关产品推荐

