Spark SQL(Scala):如何计算作者最老书籍的最小章节页数
问题:计算作者最老书籍的最小章节页数
我们有一张作者表authorsDF,其中books字段为Seq[Book]类型,每本书的chapters字段为Seq[Chapter]类型。需求是生成一张新表,每行对应一位作者,新增minChapterPage列,代表该作者最老书籍(即releaseTimestamp最小的书籍)中的最小章节页数。
示例数据定义
case class Chapter(chapterTitle: String, pages: Int) case class Book(title: String, releaseTimestamp: BigInt, chapters: Seq[Chapter]) case class Author(id: Int, name: String, nationality: String, books: Seq[Book]) val chapter1 = Chapter(chapterTitle="A", pages=23) val chapter2 = Chapter(chapterTitle="B", pages=31) val chapter3 = Chapter(chapterTitle="C", pages=51) val chapter4 = Chapter(chapterTitle="D", pages=178) val chapter5 = Chapter(chapterTitle="E", pages=12) val chapter6 = Chapter(chapterTitle="F", pages=23) val chapter7 = Chapter(chapterTitle="G", pages=4) val chapter8 = Chapter(chapterTitle="H", pages=46) val chapter9 = Chapter(chapterTitle="I", pages=30) val book1 = Book(title="Harry Potter", releaseTimestamp=1023131, chapters=Seq(chapter1, chapter2)) val book2 = Book(title="Fantastic Beasts", releaseTimestamp=1514322, chapters=Seq(chapter3)) val book3 = Book(title="Mistborn", releaseTimestamp=172322, chapters=Seq(chapter4, chapter5)) val book4 = Book(title="The Way of Kings", releaseTimestamp=651231, chapters=Seq(chapter6, chapter7)) val book5 = Book(title="A Game of Thrones", releaseTimestamp=812312, chapters=Seq(chapter8, chapter9)) val author1 = Author(id=1, name="J K Rowling", nationality="UK", books=Seq(book1, book2)) val author2 = Author(id=2, name="Brandon Sanderson", nationality="US", books=Seq(book3, book4)) val author3 = Author(id=3, name="George R R Martin", nationality="US", books=Seq(book5)) val table = Seq(author1, author2, author3) val authorsDF = table.toDF()
初始表结构
| id | name | nationality | books |
|---|---|---|---|
| 1 | J K Rowling | UK | 书籍数组内容... |
| 2 | Brandon Sanderson | US | 书籍数组内容... |
| 3 | George R R Martin | US | 书籍数组内容... |
期望输出结果
| id | name | nationality | minChapterPage |
|---|---|---|---|
| 1 | J K Rowling | UK | 23 |
| 2 | Brandon Sanderson | US | 12 |
| 3 | George R R Martin | US | 30 |
解决方案
方法一:展开+聚合(分步清晰)
通过展开数组、筛选最老书籍、再聚合最小章节页数的方式实现,步骤直观:
import org.apache.spark.sql.functions._ // 1. 展开books数组,得到每个作者的单本书记录 val authorsWithBooksDF = authorsDF.select( col("id"), col("name"), col("nationality"), explode(col("books")).alias("book") ) // 2. 计算每个作者最老书籍的发布时间戳 val authorOldestTimestampDF = authorsWithBooksDF.groupBy("id") .agg(min(col("book.releaseTimestamp")).alias("oldest_timestamp")) // 3. 筛选出每个作者的最老书籍 val authorOldestBookDF = authorsWithBooksDF.join( authorOldestTimestampDF, col("id") === authorOldestTimestampDF("id") && col("book.releaseTimestamp") === authorOldestTimestampDF("oldest_timestamp"), "inner" ).drop(authorOldestTimestampDF("id"), "oldest_timestamp") // 4. 展开chapters数组,得到章节页数 val authorChaptersDF = authorOldestBookDF.select( col("id"), col("name"), col("nationality"), explode(col("book.chapters")).alias("chapter") ) // 5. 按作者聚合,取最小章节页数 val resultDF = authorChaptersDF.groupBy("id", "name", "nationality") .agg(min(col("chapter.pages")).alias("minChapterPage")) // 查看结果 resultDF.show()
方法二:高阶函数(简洁高效)
利用Spark内置高阶函数,无需多次展开和关联,直接在每行内处理数组,效率更高:
import org.apache.spark.sql.functions._ val resultDF = authorsDF.withColumn( "minChapterPage", expr(""" array_min( transform( filter(books, b -> b.releaseTimestamp = array_min(books.releaseTimestamp)), b -> array_min(b.chapters.pages) ) ) """) ).drop("books") // 查看结果 resultDF.show()
表达式说明:
array_min(books.releaseTimestamp):获取当前作者所有书籍的最小发布时间戳(最老书籍的时间)filter(books, b -> b.releaseTimestamp = ...):筛选出该作者的所有最老书籍(兼容多本同时间的情况)transform(..., b -> array_min(b.chapters.pages)):对每本筛选出的书籍,计算其章节中的最小页数array_min(...):如果有多本最老书籍,取这些书籍的最小章节页数中的最小值
内容的提问来源于stack exchange,提问作者Paul Eaten
相关产品推荐
相关产品推荐

