Spark Structured Streaming读取CSV报basePath必须为目录错误求助
解决Spark Structured Streaming中 "Option 'basePath' must be a directory" 报错
嘿,我来帮你搞定这个Structured Streaming的报错问题!你在Spark 2.2.1里尝试读取单个CSV文件做流处理时触发了这个异常,咱们一步步来分析解决:
问题根源
先看你代码里的读取部分:
val ds = spark.readStream.schema(schema).format("csv").option("header","false").option("sep", ";").load("file:///tmp/kafka-sample-messages.csv")
这里你直接指定了单个CSV文件的路径,但Spark Structured Streaming的文件数据源(FileStreamSource)是专门用来监听目录的,它会持续扫描目录下新增的文件。当你传入单个文件路径时,Spark会默认把它当作目录来处理,这就导致了报错里的Option 'basePath' must be a directory异常。
快速解决方法
只需要把读取路径改成文件所在的目录路径就行,具体步骤如下:
- 创建一个目录(比如
/tmp/kafka-messages-dir/),把你的kafka-sample-messages.csv文件移动到这个目录里 - 修改代码中的
load路径为目录路径
修改后的完整代码
import org.apache.spark.sql.Encoders import scala.concurrent.duration._ import org.apache.spark.sql.streaming.{OutputMode, Trigger} sc.setLogLevel("INFO") case class KafkaMessage(topic: String, id: String, data: String) val schema = Encoders.product[KafkaMessage].schema // 替换为目录路径 val ds = spark.readStream.schema(schema).format("csv").option("header","false").option("sep", ";").load("file:///tmp/kafka-messages-dir/").as[KafkaMessage] val msgs = ds.groupBy('id).agg(count('id) as "total") val msgsStream = msgs.writeStream.format("console").outputMode(OutputMode.Complete).queryName("textStream").start
补充说明
- 运行修改后的代码后,Spark会先处理目录中已有的
kafka-sample-messages.csv文件,输出聚合结果到控制台 - 后续如果有新的CSV文件添加到这个目录,流任务会自动检测到并处理这些新文件,完美符合流处理的需求
- 如果只是想测试单个文件的处理逻辑,这个方法完全适用,不需要额外配置
内容的提问来源于stack exchange,提问作者Kleyson Rios
相关产品推荐
相关产品推荐

