Spark Scala文件流读取文件夹失败求助:代码修正(Windows7环境)
修正后的Spark Streaming文件流读取(Windows7环境)
我来帮你搞定这个问题!刚接触Spark和Scala遇到路径和文件流的坑很正常,尤其是Windows环境下确实有一些需要注意的细节。下面是修正后的完整代码,以及关键要点的说明:
完整修正代码
import org.apache.spark.SparkConf import org.apache.spark.sql.SparkSession import org.apache.spark.streaming.{Seconds, StreamingContext} import com.datastax.spark.connector.streaming._ object FileStreamToCassandra { def main(args: Array[String]): Unit = { // 1. 配置SparkSession,添加Cassandra连接信息 val sparkConf = new SparkConf() .setAppName("FileStreamToCassandra") .setMaster("local[*]") // 本地测试用,生产环境去掉 .set("spark.cassandra.connection.host", "你的Cassandra地址") // 比如localhost .set("spark.cassandra.connection.port", "9042") val spark = SparkSession.builder().config(sparkConf).getOrCreate() val sc = spark.sparkContext sc.setLogLevel("WARN") // 减少日志输出,方便查看关键信息 // 2. 创建StreamingContext,批次间隔1秒 val ssc = new StreamingContext(sc, Seconds(1)) // 3. 修正文件路径读取(Windows下的正确格式) // 注意:textFileStream只能监控新增文件,已存在的文件不会被读取 val lines = ssc.textFileStream("file:///C:/input/") // 4. 将流数据持久化到Cassandra // 假设你的Cassandra表结构是:keyspace.table_name (id int, content text) lines.foreachRDD { rdd => if (!rdd.isEmpty()) { // 将RDD转换为DataFrame或者直接用Cassandra连接器写入 import spark.implicits._ val df = rdd.zipWithIndex().map { case (line, idx) => (idx.toInt, line) }.toDF("id", "content") df.write .format("org.apache.spark.sql.cassandra") .options(Map("keyspace" -> "你的keyspace名称", "table" -> "你的表名称")) .mode("append") .save() } } // 5. 启动流任务并等待结束 ssc.start() ssc.awaitTermination() } }
关键修改与注意事项
- 路径格式修正:Windows下
file:///C:/input/是正确的格式,注意三个斜杠开头,后面跟盘符和路径。也可以用本地路径写法"C:/input/",但推荐用file:///的URI格式避免歧义。 - textFileStream的限制:这个API只能监控新增到目录的文件,程序启动前已经存在的文件不会被读取。如果需要读取历史文件,建议用批处理+流处理结合的方式,或者改用
readStream(Structured Streaming,更推荐的新版本API)。 - 文件写入要求:往
C:/input/目录添加文件时,必须是原子性写入(比如先在其他目录写完文件,再移动到监控目录),如果直接在监控目录下修改或创建文件,Spark可能无法识别。 - Cassandra依赖:确保你的项目中添加了Spark-Cassandra连接器的依赖,比如Maven依赖:
<dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector_2.12</artifactId> <version>3.4.1</version> <!-- 对应你的Spark版本 --> </dependency> - 本地运行注意:本地测试时
setMaster("local[*]")要设置,因为Spark Streaming需要至少2个线程(一个监控文件,一个处理数据),如果用local[1]会导致任务阻塞。
更推荐的Structured Streaming版本(替代DStream)
如果你的Spark版本是2.0+,更推荐用Structured Streaming,它的API更简洁,也更稳定,文件读取的支持更好:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.streaming.Trigger object StructuredFileStreamToCassandra { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("StructuredFileStreamToCassandra") .master("local[*]") .config("spark.cassandra.connection.host", "localhost") .getOrCreate() import spark.implicits._ // 读取目录下的文本文件流 val streamDF = spark.readStream .text("C:/input/") // 写入Cassandra val query = streamDF.writeStream .foreachBatch { (batchDF, _) => if (!batchDF.isEmpty) { batchDF.zipWithIndex().map { case (line, idx) => (idx.toInt, line.getString(0)) } .toDF("id", "content") .write .format("org.apache.spark.sql.cassandra") .options(Map("keyspace" -> "your_keyspace", "table" -> "your_table")) .mode("append") .save() } } .trigger(Trigger.ProcessingTime("1 second")) .start() query.awaitTermination() } }
内容的提问来源于stack exchange,提问作者Gnana
相关产品推荐
相关产品推荐

