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

Spark中无需保存即可将CSV转换为Parquet的方法求助

无需保存文件直接使用CSV转Parquet格式数据的实现方式

你不需要先将CSV转换后的Parquet数据写入磁盘,Spark读取CSV后生成的DataFrame本身就是可以直接用于后续操作的分布式数据集,和从Parquet文件读取的DataFrame在使用上没有任何区别。之前的写入操作只是用于持久化数据到磁盘,若仅需在Spark任务内部继续处理,完全可以省略这一步。

Python 示例代码

import sys
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("CSVProcessing").getOrCreate()

# 读取CSV文件到DataFrame
data = spark.read.load(
    sys.argv[1],
    format="csv",
    sep="|",
    inferSchema="true",
    header="true"
)

# 直接基于DataFrame执行后续操作(示例:过滤+聚合)
filtered_data = data.filter(data["age"] > 30)
aggregated_result = filtered_data.groupBy("department").count()

# 输出结果或继续其他操作
aggregated_result.show()

Scala 示例代码

import org.apache.spark.sql.SparkSession

object CSVInMemoryProcessing {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder.appName("CSVProcessing").getOrCreate()
    
    // 读取CSV文件到DataFrame
    val data = spark.read
      .format("csv")
      .option("sep", "|")
      .option("inferSchema", "true")
      .option("header", "true")
      .load(args(0))
    
    // 直接执行后续操作(示例:过滤+聚合)
    val filteredData = data.filter(data("age") > 30)
    val aggregatedResult = filteredData.groupBy("department").count()
    
    // 输出结果或继续其他操作
    aggregatedResult.show()
  }
}

优化建议

如果后续操作需要多次使用该数据集,可以通过cache()或persist()方法将DataFrame缓存到内存(或磁盘),避免重复读取CSV文件带来的性能损耗:

Python 缓存示例

data = spark.read.load(
    sys.argv[1],
    format="csv",
    sep="|",
    inferSchema="true",
    header="true"
).cache()

Scala 缓存示例

val data = spark.read
  .format("csv")
  .option("sep", "|")
  .option("inferSchema", "true")
  .option("header", "true")
  .load(args(0))
  .cache()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 05:50:41