【发布时间】:2021-10-18 22:49:18
【问题描述】:
基于特定键聚合数据集,但将聚合列表限制为固定数量。
附上创建数据集的代码。需要帮助来实现类似于 grouped() 与列表一起使用的机制。
case class AggrBook(
city: String,
state:String,
books:List[Int]
)
case class Bookings(bookingId: Int,
userId:String,
city: String,
state:String
)
val spark = SparkSession.builder.master("local")getOrCreate()
import spark.sqlContext.implicits._
val bookDS = spark.createDataset (
Seq(
Bookings(1, "ames", "Eureka", "CA"),
Bookings(2, "cha", "Eureka", "CA"),
Bookings(3, "ygy", "Eureka", "CA"),
Bookings(4, "6fsj", "Kettlemen", "AK"),
Bookings(5, "skj", "Eureka", "CA"),
Bookings(6, "po", "Irvine", "CA")
)
)
bookDS.show
val dsGrouped: Dataset[AggrBook] = bookDS.groupByKey(r => (r.city, r.state))
.mapGroups((key, value) => AggrBook(key._2, key._1, value.map(_.bookingId).toList))
dsGrouped.show()
按最多 2 个 bookingID 分组,按州、每条记录的城市聚合数据集。
我的结果
+----+---------+------------+
|city| state| books|
+----+---------+------------+
| CA| Eureka|[1, 2, 3, 5]|
| AK|Kettlemen| [4]|
| CA| Irvine| [6]|
+----+---------+------------+
期待:
+----+---------+------------+
|city| state| books|
+----+---------+------------+
| CA| Eureka|[1, 2] |
| CA| Eureka|[3, 5] |
| AK|Kettlemen|[4] |
| CA| Irvine|[6] |
+----+---------+------------+
【问题讨论】:
-
您能否更新您的问题以包含
myDataSet的定义?包括它的类型定义。例如DataFrame或Dataset[Something]。 -
已更新。谢谢。
-
你能把这行加进去吗:
val myDataSet: <type-here> = ... -
用示例数据更新了问题
标签: scala apache-spark-dataset