Spark窗口聚合结果CSV写入方法及写入报错问题排查
问题解答
1. 如何以CSV格式写入窗口聚合结果?
要把窗口聚合结果写入CSV,核心得先处理窗口函数生成的Struct类型列——因为CSV数据源不支持直接存储struct、array这类复杂数据类型。具体可以按以下步骤操作:
拆分窗口结构体列
窗口函数返回的window列是包含start和end的结构体,我们需要把它拆成两个独立的普通列,比如命名为window_start和window_end,同时移除原有的struct列:val df_agg_processed = df_agg_without_time .withColumn("window_start", $"window.start") .withColumn("window_end", $"window.end") .drop("window")正常写入CSV
处理完列之后,就可以用writeStream配置CSV写入了,记得设置必要的参数:df_agg_processed .writeStream .outputMode("append") .partitionBy("xml_data_dt") .format("csv") .option("header", "true") // 可选:写入表头,方便后续查看数据 .option("path", "hdfs://op/apps/hive/warehouse/area.db/finalTable_repo") .option("checkpointLocation", "/user/sas/sparkCheckpoint/csv_write") // 建议单独设置该查询的checkpoint路径 .trigger(Trigger.ProcessingTime("2 seconds")) .start()
2. 错误原因与解决办法
根本原因
你遇到的UnsupportedOperationException: CSV data source does not support struct<start:timestamp,end:timestamp> data type错误,本质很简单:
Spark的CSV数据源不支持写入复杂数据类型(比如struct、array、map),而你的聚合代码中groupBy(window(...))生成了一个名为window的struct类型列(包含start和end两个时间戳字段),直接写入CSV时就触发了这个不支持的报错。
针对你代码的修改方案
只需要在聚合之后、写入之前,把window结构体拆成单独列即可,修改后的完整代码片段如下:
val spark = SparkSession .builder .enableHiveSupport() .config("hive.exec.dynamic.partition", "true") .config("hive.exec.dynamic.partition.mode", "nonstrict") .config("spark.sql.streaming.checkpointLocation", "/user/sas/sparkCheckpoint") .getOrCreate // 你的原聚合逻辑 val df_agg_without_time = sqlResultjoin .withWatermark("event_time", "10 seconds") .groupBy( window($"event_time", "10 seconds", "5 seconds"), $"section", $"timestamp") .agg(sum($"total") as "total") // 新增:拆分window结构体,移除原struct列 val df_agg_for_csv = df_agg_without_time .withColumn("window_start", $"window.start") .withColumn("window_end", $"window.end") .drop("window") // 用处理后的DataFrame写入CSV df_agg_for_csv .writeStream .outputMode("append") .partitionBy("xml_data_dt") .format("csv") .option("header", "true") // 可选,建议加上 .trigger(Trigger.ProcessingTime("2 seconds")) .option("path", "hdfs://op/apps/hive/warehouse/area.db/finalTable_repo") .start()
这样修改后,DataFrame里就没有struct类型的列了,CSV写入就能正常执行。
内容的提问来源于stack exchange,提问作者BigD
相关产品推荐
相关产品推荐

