Pra*_*ade 2 streaming apache-spark spark-streaming apache-spark-sql csv-write-stream
我正在使用火花流从 Kafka 读取数据并传递到 py 文件进行预测。它返回预测以及原始数据。它将原始数据及其预测保存到文件中,但是它为每个 RDD 创建了一个文件。我需要一个由收集到的所有数据组成的单个文件,直到我停止将程序保存到单个文件中。
我试过 writeStream 它甚至不创建单个文件。我尝试使用 append 将它保存到 parquet,但它创建了多个文件,每个 RDD 为 1。我试图用追加模式写多个文件作为输出。下面的代码创建一个文件夹 output.csv 并将所有文件输入其中。
def main(args: Array[String]): Unit = {
val ss = SparkSession.builder()
.appName("consumer")
.master("local[*]")
.getOrCreate()
val scc = new StreamingContext(ss.sparkContext, Seconds(2))
val kafkaParams = Map[String, Object](
"bootstrap.servers" -> "localhost:9092",
"key.deserializer"->
"org.apache.kafka.common.serialization.StringDeserializer",
"value.deserializer">
"org.apache.kafka.common.serialization.StringDeserializer",
"group.id"-> "group5" // clients can take
)
mappedData.foreachRDD(
x =>
x.map(y =>
ss.sparkContext.makeRDD(List(y)).pipe(pyPath).toDF().repartition(1)
.write.format("csv").mode("append").option("truncate","false")
.save("output.csv")
)
)
scc.start()
scc.awaitTermination()
Run Code Online (Sandbox Code Playgroud)
我只需要获取 1 个文件,其中包含流式传输时收集的所有语句。
任何帮助将不胜感激,谢谢您的期待。
写入后,您无法修改 hdfs 中的任何文件。如果您希望实时写入文件(每 2 秒将来自流作业的数据块附加到同一文件中),则不允许这样做,因为 hdfs 文件是不可变的。如果可能,我建议您尝试编写一个读取多个文件的读取逻辑。
但是,如果您必须从单个文件中读取,我建议在您将输出写入单个 csv/parquet 文件夹后,使用“Append” SaveMode(它将为您每次写入的每个块创建部分文件)中的任何一种2 秒)。
您可以在 spark 中编写一个简单的逻辑来读取包含多个文件的此文件夹,并使用 reparation(1) 或 coalesce(1) 将其作为单个文件写入另一个 hdfs 位置,然后从该位置读取数据。见下文:
spark.read.csv("oldLocation").coalesce(1).write.csv("newLocation")
Run Code Online (Sandbox Code Playgroud)| 归档时间: |
|
| 查看次数: |
2910 次 |
| 最近记录: |