Pun*_*Raj 5 apache-spark apache-spark-sql spark-dataframe apache-spark-dataset apache-spark-xml
我正在尝试从JavaRDd <Book>和JavaRdd <Reviews>生成一个复杂的xml,我如何将这两个结合在一起以在xml之下生成?
<xml>
<library>
<books>
<book>
<author>test</author>
</book>
</books>
<reviews>
<review>
<id>1</id>
</review>
</reviews>
</library>
Run Code Online (Sandbox Code Playgroud)
如您所见,有一个父根库,其中包含子书和评论。
以下是我如何生成Book and Review Dataframe
DataFrame bookFrame = sqlCon.createDataFrame(bookRDD, Book.class);
DataFrame reviewFrame = sqlCon.createDataFrame(reviewRDD, Review.class);
Run Code Online (Sandbox Code Playgroud)
我知道要生成xml,而我的疑问尤其是对于拥有Library rootTag以及将Books and Reviews作为其子元素。
我正在使用Java。但是如果您可以指出正确的内容,则可以编写Scala或Python示例。
它可能不是使用 Spark 执行此操作的最有效方法,但下面的代码可以按照您的意愿工作。(不过是 Scala,因为我的 Java 有点生疏了)
import java.io.{File, PrintWriter}
import org.apache.spark.sql.{SaveMode, SparkSession}
import scala.io.Source
val spark = SparkSession.builder()
.master("local[3]")
.appName("test")
.config("spark.driver.allowMultipleContexts", "true")
.getOrCreate()
import spark.implicits._
/* Some code to test */
case class Book(author: String)
case class Review(id: Int)
org.apache.spark.sql.catalyst.encoders.OuterScopes.addOuterScope(this)
val bookFrame = List(
Book("book1"),
Book("book2"),
Book("book3"),
Book("book4"),
Book("book5")
).toDS()
val reviewFrame = List(
Review(1),
Review(2),
Review(3),
Review(4)
).toDS()
/* End test code **/
// Using databricks api save as 1 big xml file (instead of many parts, using repartition)
// You don't have to use repartition, but each part-xxx file will wrap contents in the root tag, making it harder to concat later.
// And TBH it really doesn't matter that Spark is doing the merging here, since the combining of data is already on the master node only
bookFrame
.repartition(1)
.write
.format("com.databricks.spark.xml")
.option("rootTag", "books")
.option("rowTag", "book")
.mode(SaveMode.Overwrite)
.save("/tmp/books/") // store to temp location
// Same for reviews
reviewFrame
.repartition(1)
.write
.format("com.databricks.spark.xml")
.option("rootTag", "reviews")
.option("rowTag", "review")
.mode(SaveMode.Overwrite)
.save("/tmp/review") // store to temp location
def concatFiles(path:String):List[String] = {
new File(path)
.listFiles
.filter(
_.getName.startsWith("part") // get all part-xxx files only (should be only 1)
)
.flatMap(file => Source.fromFile(file.getAbsolutePath).getLines())
.map(" " + _) // prefix with spaces to allow for new root level xml
.toList
}
val lines = List("<xml>","<library>") ++ concatFiles("/tmp/books/") ++ concatFiles("/tmp/review/") ++ List("</library>")
new PrintWriter("/tmp/target.xml"){
write(lines.mkString("\n"))
close
}
Run Code Online (Sandbox Code Playgroud)
结果:
<xml>
<library>
<books>
<book>
<author>book1</author>
</book>
<book>
<author>book2</author>
</book>
<book>
<author>book3</author>
</book>
<book>
<author>book4</author>
</book>
<book>
<author>book5</author>
</book>
</books>
<reviews>
<review>
<id>1</id>
</review>
<review>
<id>2</id>
</review>
<review>
<id>3</id>
</review>
<review>
<id>4</id>
</review>
</reviews>
</library>
Run Code Online (Sandbox Code Playgroud)
另一种方法可能是(仅使用 Spark)创建一个新对象
case class BookReview(books: List[Book], reviews: List[Review])并将其存储到 xml 中,然后将.collect()所有书籍和评论存储到单个列表中。
虽然我不会使用 Spark 只处理单个记录(BookReview),而是使用普通的 xml 库(如 xstream 等)来存储该对象。
更新列表 concat 方法对内存不友好,因此使用流和缓冲区这可能是一种解决方案而不是该concatFiles方法。
def outputConcatFiles(path: String, outputFile: File): Unit = {
new File(path)
.listFiles
.filter(
_.getName.startsWith("part") // get all part-xxx files only (should be only 1)
)
.foreach(file => {
val writer = new BufferedOutputStream(new FileOutputStream(outputFile, true))
val reader = new BufferedReader(new InputStreamReader(new FileInputStream(file)))
try {
Stream.continually(reader.readLine())
.takeWhile(_ != null)
.foreach(line =>
writer.write(s" $line\n".getBytes)
)
} catch {
case e: Exception => println(e.getMessage)
} finally {
writer.close()
reader.close()
}
})
}
val outputFile = new File("/tmp/target2.xml")
new PrintWriter(outputFile) { write("<xml>\n<library>\n"); close}
outputConcatFiles("/tmp/books/", outputFile)
outputConcatFiles("/tmp/review/", outputFile)
new PrintWriter(new FileOutputStream(outputFile, true)) { append("</library>"); close}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
1252 次 |
| 最近记录: |