Spark中如何按年计算平均收盘价?StringType日期列转换受阻
解决Spark SQL中字符串日期转年份计算平均收盘价的问题
咱们先明确问题根源:你的dt字段是StringType,但Spark内置的year()函数要求输入是DateType或TimestampType,所以直接调用year(dt)会失败。下面给你三种可行的解决办法,你可以根据场景选择:
方法一:在SQL查询里直接转换日期格式
这种方法不用修改原DataFrame,直接在查询语句中用to_date()函数把字符串转成日期类型,再提取年份。需要注意必须指定和你的数据匹配的日期格式(比如常见的yyyy-MM-dd、MM/dd/yyyy等):
假设你的dt字段格式是yyyy-MM-dd(比如2023-10-05),修改后的查询语句如下:
sqlContext.sql(""" SELECT AVG(closeprice), YEAR(to_date(dt, 'yyyy-MM-dd')) AS year FROM df0 GROUP BY YEAR(to_date(dt, 'yyyy-MM-dd')) """).show()
如果你的日期格式是其他样式(比如MM/dd/yyyy),把to_date的第二个参数改成对应的格式字符串即可。
方法二:先给DataFrame添加日期类型列再查询
如果需要多次使用日期字段,建议先给DataFrame新增一个日期类型的列,后续查询直接用这个新列更方便:
import org.apache.spark.sql.functions._ // 新增dt_date列,将原字符串转成日期类型 val dfWithDate = df0.withColumn("dt_date", to_date(col("dt"), "yyyy-MM-dd")) // 用新的日期列执行查询 sqlContext.sql(""" SELECT AVG(closeprice), YEAR(dt_date) AS year FROM dfWithDate GROUP BY YEAR(dt_date) """).show()
方法三:读取数据时直接指定dt为日期类型
从源头解决问题,在读取CSV的时候就把dt字段定义为DateType,并指定日期格式,这样后续查询就可以直接用原字段:
import org.apache.spark.sql.types.{StructType, StructField, DateType, DoubleType, IntegerType} val customSchema = StructType(Array( StructField("dt", DateType, true), // 将原StringType改为DateType StructField("openprice", DoubleType, true), StructField("highprice", DoubleType, true), StructField("lowprice", DoubleType, true), StructField("closeprice", DoubleType, true), StructField("volume", IntegerType, true) )) val df0 = sqlContext.read .format("com.databricks.spark.csv") .option("delimiter", ",") .option("header", "true") .option("dateFormat", "yyyy-MM-dd") // 匹配你的日期格式 .schema(customSchema) .load("./data/GSPC.csv") // 现在原查询语句可以直接运行了 sqlContext.sql("SELECT avg(closeprice), year(dt) FROM df0 GROUP BY year(dt)").show()
注意事项
- 一定要确保
to_date或dateFormat指定的格式和你数据中dt字段的实际格式完全一致,否则转换会得到null值,导致计算结果不准确。 - 如果你的日期包含时间部分(比如
2023-10-05 14:30:00),可以用to_timestamp()函数代替to_date(),后续提取年份的逻辑是一样的。
内容的提问来源于stack exchange,提问作者madeinQuant
相关产品推荐
相关产品推荐

