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

