{"id":381998,"date":"2024-06-29T04:00:24","date_gmt":"2024-06-29T04:00:24","guid":{"rendered":"http:\/\/savepearlharbor.com\/?p=381998"},"modified":"-0001-11-30T00:00:00","modified_gmt":"-0001-11-29T21:00:00","slug":"","status":"publish","type":"post","link":"https:\/\/savepearlharbor.com\/?p=381998","title":{"rendered":"<span>\u0418\u0437\u043c\u0435\u043d\u0438\u0442\u044c \u0441\u043e\u0445\u0440\u0430\u043d\u0435\u043d\u0438\u044f Spark \u0427\u0430\u0441\u0442\u044c \u0432\u0442\u043e\u0440\u0430\u044f: \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044f \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440\u0430<\/span>"},"content":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p><strong><em>\u0410\u0432\u0442\u043e\u0440:<\/em><\/strong><em> \u0418\u0432\u0430\u043d \u041a\u0430\u043b\u0438\u043d\u0438\u043d\u0441\u043a\u0438\u0439, \u0443\u0447\u0430\u0441\u0442\u043d\u0438\u043a \u043f\u0440\u043e\u0444\u0435\u0441\u0441\u0438\u043e\u043d\u0430\u043b\u044c\u043d\u043e\u0433\u043e \u0441\u043e\u043e\u0431\u0449\u0435\u0441\u0442\u0432\u0430 \u0421\u0431\u0435\u0440\u0430 SberProfi DWH\/BigData.<\/em><\/p>\n<p><em>\u041f\u0440\u043e\u0444\u0435\u0441\u0441\u0438\u043e\u043d\u0430\u043b\u044c\u043d\u043e\u0435 \u0441\u043e\u043e\u0431\u0449\u0435\u0441\u0442\u0432\u043e SberProfi DWH\/BigData \u043e\u0442\u0432\u0435\u0447\u0430\u0435\u0442 \u0437\u0430 \u0440\u0430\u0437\u0432\u0438\u0442\u0438\u0435 \u043a\u043e\u043c\u043f\u0435\u0442\u0435\u043d\u0446\u0438\u0439 \u0432 \u0442\u0430\u043a\u0438\u0445 \u043d\u0430\u043f\u0440\u0430\u0432\u043b\u0435\u043d\u0438\u044f\u0445, \u043a\u0430\u043a \u044d\u043a\u043e\u0441\u0438\u0441\u0442\u0435\u043c\u0430 Hadoop, Teradata, Oracle DB, GreenPlum, \u0430 \u0442\u0430\u043a\u0436\u0435 BI \u0438\u043d\u0441\u0442\u0440\u0443\u043c\u0435\u043d\u0442\u0430\u0445 Qlik, SAP BO, Tableau \u0438 \u0434\u0440.<\/em><\/p>\n<p>\u041d\u0430\u0447\u043d\u0443 \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u0435 \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440\u0430 \u0441 UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u044b:<\/p>\n<figure class=\"full-width\"><img loading=\"lazy\" decoding=\"async\" src=\"https:\/\/habrastorage.org\/r\/w1560\/getpro\/habr\/upload_files\/964\/e1d\/99d\/964e1d99d9eebd2740a09ba7ca1be2dd.png\" alt=\"\u0420\u0438\u0441\u0443\u043d\u043e\u043a 1. UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u0430 \u043a\u043b\u0430\u0441\u0441\u043e\u0432 OrderBucketsPartitioner\" title=\"\u0420\u0438\u0441\u0443\u043d\u043e\u043a 1. UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u0430 \u043a\u043b\u0430\u0441\u0441\u043e\u0432 OrderBucketsPartitioner\" width=\"821\" height=\"1061\" data-src=\"https:\/\/habrastorage.org\/getpro\/habr\/upload_files\/964\/e1d\/99d\/964e1d99d9eebd2740a09ba7ca1be2dd.png\"\/><figcaption>\u0420\u0438\u0441\u0443\u043d\u043e\u043a 1. UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u0430 \u043a\u043b\u0430\u0441\u0441\u043e\u0432 OrderBucketsPartitioner<\/figcaption><\/figure>\n<p>\u041d\u0430\u0432\u0435\u0440\u043d\u044f\u043a\u0430 \u0432\u044b \u0432\u0438\u0434\u0435\u043b\u0438, \u043a\u0430\u043a Spark \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0435\u0442 \u0437\u0430\u043f\u0440\u043e\u0441, \u0441\u0442\u0440\u043e\u0438\u0442 \u043f\u043b\u0430\u043d \u0438 \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442 \u0435\u0433\u043e. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u0438\u043c \u044d\u0442\u0443 \u0441\u0445\u0435\u043c\u0443 \u0441 \u043d\u0430\u0448\u0435\u0439 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0435\u0439 (\u043a\u043e\u0433\u0434\u0430 \u0431\u0443\u0434\u0435\u0442\u0435 \u0437\u043d\u0430\u043a\u043e\u043c\u0438\u0442\u044c\u0441\u044f \u0441 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0435\u0439 \u043a\u043b\u0430\u0441\u0441\u043e\u0432, \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0439\u0442\u0435\u0441\u044c \u043a \u044d\u0442\u043e\u0439 \u0441\u0445\u0435\u043c\u0435, \u0447\u0442\u043e\u0431\u044b \u0443\u0432\u0438\u0434\u0435\u0442\u044c \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0438\u0435):<\/p>\n<figure class=\"full-width\"><img loading=\"lazy\" decoding=\"async\" src=\"https:\/\/habrastorage.org\/r\/w1560\/getpro\/habr\/upload_files\/eab\/499\/e37\/eab499e37f38678093879d9be9b058f1.png\" alt=\"\u0420\u0438\u0441\u0443\u043d\u043e\u043a 2. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0445\u0435\u043c\u044b \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0437\u0430\u043f\u0440\u043e\u0441\u0430 \u0438 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438.\" title=\"\u0420\u0438\u0441\u0443\u043d\u043e\u043a 2. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0445\u0435\u043c\u044b \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0437\u0430\u043f\u0440\u043e\u0441\u0430 \u0438 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438.\" width=\"1082\" height=\"192\" data-src=\"https:\/\/habrastorage.org\/getpro\/habr\/upload_files\/eab\/499\/e37\/eab499e37f38678093879d9be9b058f1.png\"\/><figcaption>\u0420\u0438\u0441\u0443\u043d\u043e\u043a 2. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0445\u0435\u043c\u044b \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0437\u0430\u043f\u0440\u043e\u0441\u0430 \u0438 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438.<\/figcaption><\/figure>\n<p>\u041d\u0443\u0436\u0435\u043d \u043c\u0435\u0442\u043e\u0434, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u0437\u0432\u043e\u043b\u0438\u043b \u0431\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u0434\u043b\u044f \u043b\u044e\u0431\u043e\u0433\u043e \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430. \u041f\u043e\u043a\u0430 \u043d\u0435 \u0431\u0443\u0434\u0435\u043c \u0441\u0432\u044f\u0437\u044b\u0432\u0430\u0442\u044c\u0441\u044f \u0441 SQL (\u0440\u0430\u0437\u0432\u0435 \u043a\u0442\u043e-\u0442\u043e \u0434\u0435\u043b\u0430\u0435\u0442 dataframe.repartition(&#8230;) \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e SQL?).<\/p>\n<p>\u0414\u043e\u0431\u0430\u0432\u0438\u043c \u0442\u0430\u043a\u043e\u0439 \u043c\u0435\u0442\u043e\u0434 \u0432 Scala Dataframe API, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f \u043f\u0430\u0442\u0442\u0435\u0440\u043d \u00abPimp my library\u00bb.<br \/> \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043d\u043e\u0432\u044b\u0439 \u043e\u0431\u044a\u0435\u043a\u0442 ru.kalininskii.orderbucketing.OrderBucketing. \u041e\u043d \u0431\u0443\u0434\u0435\u0442 \u0441\u043e\u0434\u0435\u0440\u0436\u0430\u0442\u044c implicit class \u0441 \u043f\u043e\u043a\u0430 \u0435\u0434\u0438\u043d\u0441\u0442\u0432\u0435\u043d\u043d\u044b\u043c \u043c\u0435\u0442\u043e\u0434\u043e\u043c repartitionWithOrderAndSort:<\/p>\n<details class=\"spoiler\">\n<summary>OrderBucketing<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing   import org.apache.spark.sql.{Column, DataFrame, PlanHelper} import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, SortOrder} import ru.kalininskii.orderbucketing.plans.logical.RepartitionWithOrderAndSort   object OrderBucketing {     implicit class DataFramePartOps(ds: DataFrame) {     def repartitionWithOrderAndSort(numLines: Int,                                      numPartitions: Int,                                      orderColumn: Column,                                      partitionColumns: Seq[Column],                                      sortColumns: Seq[Column]): DataFrame = {       def toSortOrder(col: Column): SortOrder = {         col.expr match {           case order: SortOrder => order.copy(direction = Ascending)           case expr: Expression => SortOrder(expr, Ascending)           case _ => throw new Exception(s\"Can not get order from $col\")         }       }         val orderExpression = toSortOrder(orderColumn)       val partitionExpressions = partitionColumns.map(_.expr)       val sortExpressions = sortColumns.map(toSortOrder)         val logicalPlan = RepartitionWithOrderAndSort(         orderExpression,         partitionExpressions,         sortExpressions,         numLines,         numPartitions,         None,         ds.queryExecution.logical)         PlanHelper.planToDF(ds.sparkSession, logicalPlan)     }   }   }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041a\u0430\u043a \u043c\u043e\u0436\u043d\u043e \u0432\u0438\u0434\u0435\u0442\u044c \u0438\u0437 \u043a\u043e\u0434\u0430 \u0432\u044b\u0448\u0435, \u043c\u0435\u0442\u043e\u0434 \u043f\u0440\u0438\u043d\u0438\u043c\u0430\u0435\u0442 \u0430\u0440\u0433\u0443\u043c\u0435\u043d\u0442\u00a0partitionColumns: Seq[Column], \u0432 \u043a\u043e\u0442\u043e\u0440\u043e\u043c \u043f\u043e\u043b\u044f \u044f\u0432\u043d\u043e\u0433\u043e \u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0433\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0435\u043d\u044b. \u0422\u0430\u043a \u0438 \u0434\u043e\u043b\u0436\u043d\u043e \u0431\u044b\u0442\u044c, \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u0435\u0442 \u0437\u0430\u043f\u0438\u0441\u0438 \u043f\u043e \u0441\u0435\u043a\u0446\u0438\u044f\u043c \u043d\u0435\u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e \u043e\u0442 \u0438\u0445 \u0432\u043d\u0435\u0448\u043d\u0435\u0433\u043e \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u044f.\u00a0<\/p>\n<p>\u041e\u0431\u044a\u0435\u043a\u0442 PlanHelper \u043f\u043e\u043a\u0430 \u0447\u0442\u043e \u0442\u0430\u043a\u0436\u0435 \u043e\u0431\u043e\u0439\u0434\u0451\u0442\u0441\u044f \u043e\u0434\u043d\u0438\u043c \u043c\u0435\u0442\u043e\u0434\u043e\u043c, planToDF. \u042d\u0442\u043e\u0442 \u043c\u0435\u0442\u043e\u0434 \u0432\u044b\u0437\u044b\u0432\u0430\u0435\u0442 \u043f\u0440\u0438\u0432\u0430\u0442\u043d\u044b\u0439 \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u043f\u0430\u043a\u0435\u0442\u0430 org.apache.spark.sql, \u043f\u043e\u044d\u0442\u043e\u043c\u0443, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u043c\u043e\u0433 \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0442\u044c\u0441\u044f, \u043e\u0431\u044a\u0435\u043a\u0442 \u0442\u043e\u0436\u0435 \u0434\u043e\u043b\u0436\u0435\u043d \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u044d\u0442\u043e\u043c \u043f\u0430\u043a\u0435\u0442\u0435:<\/p>\n<details class=\"spoiler\">\n<summary>PlanHelper<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package org.apache.spark.sql   import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan   object PlanHelper {     def planToDF(spark: SparkSession, logicalPlan: LogicalPlan): DataFrame = {     Dataset.ofRows(spark, logicalPlan)   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0438\u043f \u0434\u043b\u044f \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD, \u0430 \u0442\u0430\u043a\u0436\u0435 \u043a\u043b\u044e\u0447\u0430 RDD, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u043d\u0430\u043c \u0432 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0451\u043d\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442, \u043f\u043e\u043c\u0435\u0441\u0442\u0438\u043c \u0432 package object, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u0431\u044b\u043b \u0434\u043e\u0441\u0442\u0443\u043f\u0435\u043d \u0432\u0441\u0435\u043c\u0443 \u043d\u0430\u0448\u0435\u043c\u0443 \u043f\u0430\u043a\u0435\u0442\u0443 \u0438 \u043d\u0435 \u0442\u043e\u043b\u044c\u043a\u043e<\/p>\n<details class=\"spoiler\">\n<summary>package object rangebucketing<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii   import org.apache.spark.sql.catalyst.InternalRow   package object rangebucketing {   type BucketsDistribution = (Int, InternalRow, Either[(InternalRow, Seq[String]), (InternalRow, InternalRow, Int)])   type OrderAndSortKey = (InternalRow, InternalRow, InternalRow) }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0435\u043f\u0435\u0440\u044c \u043d\u0443\u0436\u043d\u043e \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u043f\u043b\u0430\u043d\u044b, \u043d\u0430\u0447\u0438\u043d\u0430\u044f \u0441 \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e. \u0418 \u0432 \u043d\u0438\u0445 \u043f\u0435\u0440\u0435\u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u043c\u0435\u0442\u043e\u0434 <em>simpleString<\/em>, \u0442\u0430\u043a \u043a\u0430\u043a \u0438\u043d\u0430\u0447\u0435 \u043e\u043d \u0432\u044b\u0432\u0435\u0434\u0435\u0442 \u0432 \u043f\u043b\u0430\u043d \u0432\u0435\u0441\u044c \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0439 \u043e\u0431\u044a\u0435\u043a\u0442 <em>distribution.\u00a0<\/em>\u041b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u0438\u0439 \u043f\u043b\u0430\u043d \u0434\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0440\u0435\u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0431\u0443\u0434\u0435\u0442 \u0432\u0438\u0434\u0435\u043d \u043f\u0440\u0438 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0438 \u043c\u0435\u0442\u043e\u0434\u0430 explain, \u0438 \u043c\u043e\u0436\u043d\u043e \u0431\u0443\u0434\u0435\u0442 \u0443\u0431\u0435\u0434\u0438\u0442\u044c\u0441\u044f, \u0447\u0442\u043e \u043e\u043d \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0435\u0442\u0441\u044f. \u041a\u0440\u043e\u043c\u0435 \u0442\u043e\u0433\u043e, \u043e\u043d \u0432\u0441\u0442\u0440\u043e\u0435\u043d \u0432 \u0444\u043e\u0440\u043c\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 \u043f\u043b\u0430\u043d\u043e\u0432 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f \u0437\u0430\u043f\u0440\u043e\u0441\u043e\u0432 \u0438 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0442\u044c \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044e \u043e \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u043e\u043c \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u043d\u0430\u0431\u043e\u0440\u0430 \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 <em>partitioning<\/em><strong><em>.<\/em><\/strong><\/p>\n<details class=\"spoiler\">\n<summary>RepartitionWithOrderAndSort<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing.plans.logical   import org.apache.spark.sql.catalyst.expressions.{Expression, SortOrder} import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, RepartitionOperation} import ru.kalininskii.orderbucketing.BucketsDistribution import ru.kalininskii.orderbucketing.plans.physical.OrderBucketsPartitioning   \/**  * \u042d\u0442\u043e\u0442 \u043a\u043b\u0430\u0441\u0441 \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u0435\u0442 \u0434\u0430\u043d\u043d\u044b\u0435 \u043f\u043e \u0443\u0441\u043b\u043e\u0432\u044e \u0443\u043d\u0438\u043a\u0430\u043b\u044c\u043d\u044b\u0445 \u0442\u043e\u0447\u043d\u044b\u0445 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 [[Expression]]  *  * @param orderExpression      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u0435, \u043f\u043e \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043a\u043e\u0442\u043e\u0440\u043e\u0433\u043e \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  *                             \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f [[partitionExpressions]]  * @param partitionExpressions \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435.  * @param sortExpressions      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u0430\u044f \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0430.  * @param numLines             \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD (\u0438 \u0432 \u0437\u0430\u043f\u0438\u0441\u0430\u043d\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435)  * @param numPartitions        \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439, \u043c\u043e\u0436\u0435\u0442 \u043e\u0442\u043b\u0438\u0447\u0430\u0442\u044c\u0441\u044f  * @param distribution         \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044f \u043e\u0431 \u0438\u043c\u0435\u044e\u0449\u0435\u043c\u0441\u044f \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0438, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043d\u0430\u0434\u043e \u0432\u043e\u0441\u043f\u0440\u043e\u0438\u0437\u0432\u0435\u0441\u0442\u0438  * @param child                \u043f\u043b\u0430\u043d \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u0438\u044f \u043d\u0430\u0431\u043e\u0440\u0430 \u0434\u0430\u043d\u043d\u044b\u0445, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043d\u0443\u0436\u043d\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u0442\u044c  *\/ case class RepartitionWithOrderAndSort(                                         orderExpression: SortOrder,                                         partitionExpressions: Seq[Expression],                                         sortExpressions: Seq[SortOrder],                                         numLines: Int,                                         numPartitions: Int,                                         distribution: Option[Seq[BucketsDistribution]],                                         child: LogicalPlan                                       ) extends RepartitionOperation {     override def nodeName: String = s\"Repartition plan with \" +     s\"${distribution.map(_ => \"predefined\").getOrElse(\"unknown\")} distribution:\"     override def simpleString: String = {     s\"$nodeName ${partitionExpressions.mkString(\"partition by [\", \", \", \"], \")}\" +       s\"global order by ${orderExpression.child}, \" +       s\"${sortExpressions.map(_.child).mkString(\"local sort by [\", \", \", \"], \")}\" +       numLines + \" rows per task, \" + numPartitions + \" initial tasks\"   }     val partitioning: OrderBucketsPartitioning = OrderBucketsPartitioning(     orderExpression, partitionExpressions, sortExpressions, numLines, numPartitions, distribution)     override def maxRows: Option[Long] = child.maxRows     override def shuffle: Boolean = true }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041a\u043b\u0430\u0441\u0441 <em>Partitioning <\/em>\u0438 \u0435\u0433\u043e \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d\u0438\u044f\u00a0\u0432 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0435 \u043f\u0430\u043a\u0435\u0442\u043e\u0432 Spark \u043e\u0442\u043d\u043e\u0441\u044f\u0442\u0441\u044f \u043a \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u0438\u043c \u043f\u043b\u0430\u043d\u0430\u043c (<em>org.apache.spark.sql.catalyst.plans.physical<\/em>), \u043d\u043e \u0441\u043a\u043e\u0440\u0435\u0435 \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u043b\u044f\u044e\u0442 \u0441\u043e\u0431\u043e\u0439 \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u044b \u0434\u043b\u044f \u0445\u0430\u0440\u0430\u043a\u0442\u0435\u0440\u0438\u0441\u0442\u0438\u043a \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0433\u043e \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430.<\/p>\n<p>\u0412\u0441\u0435 \u0432\u0441\u0442\u0440\u043e\u0435\u043d\u043d\u044b\u0435 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0441\u043e\u0437\u0434\u0430\u044e\u0442\u0441\u044f \u043a\u0435\u0439\u0441 \u043a\u043b\u0430\u0441\u0441\u043e\u043c (\u0441\u043c. \u0422\u0435\u0440\u043c\u0438\u043d\u044b \u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f)\u00a0<em>org.apache.spark.sql.catalyst.plans.logical.RepartitionByExpression,\u00a0<\/em>\u043f\u043e\u043b\u0443\u0447\u0430\u044f \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043f\u043e\u043b\u0435\u0439, \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0439 \u043a\u0435\u0439\u0441 \u043a\u043b\u0430\u0441\u0441.<\/p>\n<p>\u0418 \u0432 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435,\u00a0<em>OrderBucketsPartitioning<\/em>\u00a0\u043c\u043e\u0436\u043d\u043e \u0441\u0447\u0438\u0442\u0430\u0442\u044c \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u0435\u043c \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u0438 \u043d\u0430\u0434 \u0442\u0430\u0431\u043b\u0438\u0446\u0435\u0439, \u043f\u0440\u043e\u0447\u0438\u0442\u0430\u043d\u043d\u043e\u0439 \u0438\u0437 \u0444\u0430\u0439\u043b\u043e\u0432\u043e\u0439 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u0438\u043b\u0438 \u043e\u0431\u044a\u0435\u043a\u0442\u043d\u043e\u0433\u043e \u0445\u0440\u0430\u043d\u0438\u043b\u0438\u0449\u0430. \u041d\u0435 \u0431\u0443\u0434\u0435\u0442 \u043e\u0448\u0438\u0431\u043a\u043e\u0439 \u043e\u0442\u043d\u0435\u0441\u0442\u0438 \u043a\u043b\u0430\u0441\u0441 \u043a \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u043c\u0443 \u043f\u043b\u0430\u043d\u0443, \u0432\u0435\u0434\u044c \u043c\u0435\u0442\u043e\u0434 <em>dataframe.explain(true)<\/em> \u0438\u043c\u0435\u043d\u043d\u043e \u0435\u0433\u043e \u0432\u044b\u0432\u043e\u0434\u0438\u0442 \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0434\u0435\u0440\u0435\u0432\u0435 \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e \u043f\u043b\u0430\u043d\u0430).<\/p>\n<p>\u0412 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0435\u0441\u0442\u044c \u043d\u0435\u043f\u0440\u0438\u044f\u0442\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442: \u0442\u0440\u0435\u0439\u0442 (\u0441\u043c. \u0422\u0435\u0440\u043c\u0438\u043d\u044b \u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f) <em>Distribution \u00a0 <\/em>\u0432 \u0438\u0441\u0445\u043e\u0434\u043d\u043e\u043c \u043a\u043e\u0434\u0435 Spark \u043e\u0431\u044a\u044f\u0432\u043b\u0435\u043d \u043a\u0430\u043a <em>sealed<\/em>, \u0430 \u0437\u043d\u0430\u0447\u0438\u0442 \u043d\u0435 \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d. \u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u0432 \u043c\u0435\u0442\u043e\u0434\u0435 <em>satisfies0<\/em> \u044f \u0431\u0443\u0434\u0443 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c <em>OrderedDistribution<\/em>, \u0445\u043e\u0442\u044f \u043e\u043d\u0430 \u043d\u0435 \u043b\u0443\u0447\u0448\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c \u043f\u043e\u0434\u0445\u043e\u0434\u0438\u0442 \u0434\u043b\u044f \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f.<\/p>\n<details class=\"spoiler\">\n<summary>OrderBucketsPartitioning<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing.plans.physical   import org.apache.spark.sql.catalyst.expressions.{Expression, SortOrder, Unevaluable} import org.apache.spark.sql.catalyst.plans.physical.{Distribution, OrderedDistribution, Partitioning} import org.apache.spark.sql.types.{DataType, IntegerType} import ru.kalininskii.orderbucketing.BucketsDistribution   \/**  * \u042d\u0442\u043e\u0442 \u043a\u043b\u0430\u0441\u0441 \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u0442 \u0434\u0430\u043d\u043d\u044b\u0435 \u0434\u043b\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043f\u043e \u0434\u0438\u0430\u043f\u0430\u0437\u043e\u043d\u0430\u043c \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439  *  * @param orderExpression      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u0435, \u043f\u043e \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043a\u043e\u0442\u043e\u0440\u043e\u0433\u043e \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  *                             \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f [[partitionExpressions]]  * @param partitionExpressions \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  * @param sortExpressions      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u0430\u044f \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0430  * @param numLines             \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD (\u0438 \u0432 \u0437\u0430\u043f\u0438\u0441\u0430\u043d\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435)  * @param numPartitions        \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439, \u043c\u043e\u0436\u0435\u0442 \u043e\u0442\u043b\u0438\u0447\u0430\u0442\u044c\u0441\u044f  * @param distribution         \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044f \u043e\u0431 \u0438\u043c\u0435\u044e\u0449\u0435\u043c\u0441\u044f \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0438, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043d\u0430\u0434\u043e \u0432\u043e\u0441\u043f\u0440\u043e\u0438\u0437\u0432\u0435\u0441\u0442\u0438  *\/ case class OrderBucketsPartitioning(                                      orderExpression: SortOrder,                                      partitionExpressions: Seq[Expression],                                      sortExpressions: Seq[SortOrder],                                      numLines: Int,                                      numPartitions: Int,                                      distribution: Option[Seq[BucketsDistribution]])     extends Expression with Partitioning with Unevaluable {     override def nodeName: String = s\"Repartition with \" +     s\"${distribution.map(_ => \"predefined\").getOrElse(\"unknown\")} distribution:\"     override def simpleString: String = {     s\"$nodeName ${partitionExpressions.mkString(\"partition by [\", \", \", \"], \")}\" +       s\"global order by ${orderExpression.child}, \" +       s\"${sortExpressions.map(_.child).mkString(\"local sort by [\", \", \", \"], \")}\" +       numLines + \" rows per task, \" + numPartitions + \" initial tasks\"   }     override def children: Seq[Expression] = partitionExpressions :+ orderExpression     override def nullable: Boolean = false     override def dataType: DataType = IntegerType     override def satisfies0(required: Distribution): Boolean = {     super.satisfies0(required) || {       required match {         case o: OrderedDistribution =>           orderExpression.semanticEquals(o.ordering.head)         case _ => false       }     }   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0414\u0430\u043b\u0435\u0435 \u0434\u0432\u0430 \u043e\u0441\u043d\u043e\u0432\u043d\u044b\u0445 \u043a\u043b\u0430\u0441\u0441\u0430, <em>org.apache.spark.sql.execution.exchange.ShuffleExchangeOrderExec <\/em>\u0438 <em>org.apache.spark.sql.partitioning.OrderBucketsPartitioner<\/em>. \u041e\u043d\u0438 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044e\u0442 \u043c\u0435\u0442\u043e\u0434\u044b \u0441 \u043e\u0433\u0440\u0430\u043d\u0438\u0447\u0435\u043d\u043d\u043e\u0439 \u0432\u0438\u0434\u0438\u043c\u043e\u0441\u0442\u044c\u044e, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u0434\u043e\u043b\u0436\u043d\u044b \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043f\u0430\u043a\u0435\u0442\u0430 org.apache.spark.sql. \u042d\u0442\u0438 \u043a\u043b\u0430\u0441\u0441\u044b \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u044e\u0442 \u0442\u0440\u0430\u043d\u0441\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u0438 RDD, \u043e\u0442\u043d\u043e\u0441\u044f\u0442\u0441\u044f \u043a \u043f\u043b\u0430\u043d\u0430\u043c \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f.<\/p>\n<p>\u041e\u0441\u043d\u043e\u0432\u043d\u043e\u0435 <em>\u043e\u0442\u043b\u0438\u0447\u0438\u0435 org.apache.spark.sql.execution.exchange.ShuffleExchangeOrderExec<\/em> \u043e\u0442 \u043e\u0440\u0438\u0433\u0438\u043d\u0430\u043b\u0430 \u0437\u0430\u043a\u043b\u044e\u0447\u0430\u0435\u0442\u0441\u044f \u0432 \u0442\u043e\u043c, \u0447\u0442\u043e \u0434\u043b\u044f \u0441\u0435\u043c\u043f\u043b\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442\u0441\u044f\u00a0RDD, \u0432 \u043a\u043e\u0442\u043e\u0440\u043e\u043c <em>MutablePair<\/em> \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442 \u043e\u0431\u0435 \u0447\u0430\u0441\u0442\u0438:<\/p>\n<ul>\n<li>\n<p>\u043f\u0435\u0440\u0432\u0443\u044e &#8212; \u0434\u043b\u044f \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043f\u043e\u043b\u0435\u0439 \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f ([P]);<\/p>\n<\/li>\n<li>\n<p>\u0432\u0442\u043e\u0440\u0443\u044e &#8212; \u0434\u043b\u044f \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f ([O]).<\/p>\n<\/li>\n<\/ul>\n<p>\u0421 \u043c\u043e\u0435\u0439 \u0442\u043e\u0447\u043a\u0438 \u0437\u0440\u0435\u043d\u0438\u044f, \u044d\u0442\u043e \u0443\u043c\u0435\u043d\u044c\u0448\u0438\u0442 \u043e\u0431\u044a\u0451\u043c \u0434\u0430\u043d\u043d\u044b\u0445 \u0434\u043b\u044f \u0434\u0430\u043b\u044c\u043d\u0435\u0439\u0448\u0435\u0439 \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0438 \u043f\u043e\u0437\u0432\u043e\u043b\u0438\u0442 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u044b\u0435 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u044b \u0434\u043b\u044f \u0434\u0440\u0443\u0433\u0438\u0445 \u043d\u0430\u0431\u043e\u0440\u043e\u0432 \u0434\u0430\u043d\u043d\u044b\u0445. \u041f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u044b\u0435 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u044b \u0432 \u043b\u044e\u0431\u043e\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u0434\u043e\u043b\u0436\u043d\u044b \u0431\u044b\u0442\u044c \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u044b \u043d\u0430 \u0434\u0440\u0430\u0439\u0432\u0435\u0440, \u0430 \u0437\u0430\u0442\u0435\u043c \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u044b \u0438\u0441\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044f\u043c (\u0442\u0430\u043a \u0440\u0430\u0431\u043e\u0442\u0430\u0435\u0442 broadcast), \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u0440\u0430\u0437\u0443\u043c\u043d\u043e \u0441 \u0441\u0430\u043c\u043e\u0433\u043e \u043d\u0430\u0447\u0430\u043b\u0430 \u0443\u043c\u0435\u043d\u044c\u0448\u0430\u0442\u044c \u0438\u0445, \u043d\u0430\u0441\u043a\u043e\u043b\u044c\u043a\u043e \u044d\u0442\u043e \u0432\u043e\u0437\u043c\u043e\u0436\u043d\u043e.<\/p>\n<details class=\"spoiler\">\n<summary>ShuffleExchangeOrderExec<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package org.apache.spark.sql.execution.exchange   import org.apache.spark._ import org.apache.spark.rdd.RDD import org.apache.spark.serializer.Serializer import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.catalyst.errors._ import org.apache.spark.sql.catalyst.expressions.codegen.LazilyGeneratedOrdering import org.apache.spark.sql.catalyst.expressions.{Attribute, UnsafeProjection} import org.apache.spark.sql.catalyst.plans.physical._ import org.apache.spark.sql.execution._ import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics} import org.apache.spark.sql.partitioning.OrderBucketsPartitioner import org.apache.spark.util.MutablePair import ru.kalininskii.orderbucketing.OrderAndSortKey import ru.kalininskii.orderbucketing.plans.physical.OrderBucketsPartitioning   import scala.reflect.ClassTag   \/**  * \u0412\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442 shuffle, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043e\u043f\u0438\u0441\u044b\u0432\u0430\u0435\u0442\u0441\u044f newPartitioning  * \u041a\u043e\u0434 \u0447\u0430\u0441\u0442\u0438\u0447\u043d\u043e \u0437\u0430\u0438\u043c\u0441\u0442\u0432\u043e\u0432\u0430\u043d \u0438\u0437 [[ShuffleExchangeExec]]  *\/ case class ShuffleExchangeOrderExec(newPartitioning: OrderBucketsPartitioning,                                     child: SparkPlan,                                     partitioner: Option[Partitioner] = None) extends Exchange {     override lazy val metrics: Map[String, SQLMetric] = Map(     \"dataSize\" -> SQLMetrics.createSizeMetric(sparkContext, \"data size\"))     override def nodeName: String = \"ExchangeWithOrder\"     override def outputPartitioning: Partitioning = newPartitioning     private val serializer: Serializer = SparkEnv.get     .serializerManager.getSerializer(implicitly[ClassTag[OrderAndSortKey]],     implicitly[ClassTag[InternalRow]])     override protected def doPrepare(): Unit = {}     \/**    * \u0412\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u0442 \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c RDD [[ShuffleDependency]], \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u043f\u0435\u0440\u0435\u043c\u0435\u0448\u0430\u0435\u0442 \u0437\u0430\u043f\u0438\u0441\u0438    * \u043f\u043e \u0441\u0445\u0435\u043c\u0435, \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u043d\u043e\u0439 \u0432 `newPartitioning`. \u041f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD, \u0441\u0432\u0437\u044f\u0437\u0430\u043d\u043d\u044b\u0435 \u0441    * \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0435\u043d\u043d\u043e\u0439 ShuffleDependency \u0431\u0443\u0434\u0443\u0442 \u0432\u0445\u043e\u0434\u043d\u044b\u043c\u0438 \u0434\u0430\u043d\u043d\u044b\u043c\u0438 \u0434\u043b\u044f shuffle.    *\/   private[exchange] def prepareShuffleDependency()   : ShuffleDependency[OrderAndSortKey, InternalRow, InternalRow] = {     val rdd = child.execute()       val part: Partitioner = partitioner       .getOrElse(ShuffleExchangeOrderExec.preparePartitioner(rdd, child.output, newPartitioning))       ShuffleExchangeOrderExec.prepareShuffleDependency(rdd, child.output, newPartitioning, serializer, part)   }     \/**    * \u0412\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u0442 [[ShuffledOrderRDD]], \u0430 \u044d\u0442\u043e \u0438 \u0435\u0441\u0442\u044c \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445 \u043f\u043e\u0441\u043b\u0435 \u043f\u0435\u0440\u0435\u043c\u0435\u0448\u0438\u0432\u0430\u043d\u0438\u044f.    * [[ShuffledOrderRDD]], \u043a\u0430\u043a \u0438 \u043f\u0440\u043e\u0447\u0438\u0435 \u043d\u0430\u0441\u043b\u0435\u0434\u043d\u0438\u043a\u0438 ShuffledRDD, \u043e\u0441\u043d\u043e\u0432\u0430\u043d \u043d\u0430 \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u043e\u0439 [[ShuffleDependency]]    *\/   private[exchange] def preparePostShuffleRDD(                                                shuffleDependency                                                : ShuffleDependency[OrderAndSortKey, InternalRow, InternalRow]                                              ): ShuffledOrderRDD = {     new ShuffledOrderRDD(shuffleDependency)   }     \/**    * ShuffledOrderRDD \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u043a\u0435\u0448\u0438\u0440\u043e\u0432\u0430\u043d \u0434\u043b\u044f \u043f\u0435\u0440\u0435\u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u044f    *\/   private var cachedShuffleRDD: ShuffledOrderRDD = _     protected override def doExecute(): RDD[InternalRow] = attachTree(this, \"execute\") {     \/\/ Returns the same ShuffleRowRDD if this plan is used by multiple plans.     if (cachedShuffleRDD == null) {       cachedShuffleRDD = preparePostShuffleRDD(prepareShuffleDependency())     }     cachedShuffleRDD   }   }     object ShuffleExchangeOrderExec {     def preparePartitioner(rdd: RDD[InternalRow],                          outputAttributes: Seq[Attribute],                          partitioning: OrderBucketsPartitioning): Partitioner = {     \/\/ Internally, ValuesAndRangePartitioner runs a job on the RDD that samples keys to compute     \/\/ partition bounds. To get accurate samples, we need to copy the mutable keys.     val rddForSampling = rdd.mapPartitionsInternal { iter =>       val projectionPart = UnsafeProjection.create(partitioning.partitionExpressions, outputAttributes)       val projectionOrder = UnsafeProjection.create(partitioning.orderExpression.child :: Nil, outputAttributes)       val mutablePair = new MutablePair[InternalRow, InternalRow]()       iter.map(row => mutablePair.update(projectionPart(row).copy(), projectionOrder(row).copy()))     }     implicit val ordering: LazilyGeneratedOrdering =       new LazilyGeneratedOrdering(Seq(partitioning.orderExpression),         outputAttributes.filter(_.references equals partitioning.orderExpression.child.references))       new OrderBucketsPartitioner(       rddForSampling,       partitioning.numLines,       partitioning.numPartitions,       partitioning.distribution)   }     private def getPartitionKeyExtractor(outputAttributes: Seq[Attribute],                                        partitioning: OrderBucketsPartitioning,                                        i: Int): InternalRow => OrderAndSortKey = {     val projectionPart = UnsafeProjection.create(partitioning.partitionExpressions, outputAttributes)     val projectionOrder = UnsafeProjection.create(partitioning.orderExpression.child :: Nil, outputAttributes)     val projectionSort = UnsafeProjection.create(partitioning.sortExpressions.map(_.child), outputAttributes)     row => (projectionPart(row), projectionOrder(row), projectionSort(row).copy())   }     \/**    * \u0412\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0435\u0442 \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u044c RDD [[ShuffleDependency]], \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u043f\u0435\u0440\u0435\u043c\u0435\u0448\u0430\u0435\u0442 \u0437\u0430\u043f\u0438\u0441\u0438    * \u043f\u043e \u0441\u0445\u0435\u043c\u0435, \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u043d\u043e\u0439 \u0432 `newPartitioning`. \u041f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD, \u0441\u0432\u0437\u044f\u0437\u0430\u043d\u043d\u044b\u0435 \u0441    * \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0435\u043d\u043d\u043e\u0439 ShuffleDependency \u0431\u0443\u0434\u0443\u0442 \u0432\u0445\u043e\u0434\u043d\u044b\u043c\u0438 \u0434\u0430\u043d\u043d\u044b\u043c\u0438 \u0434\u043b\u044f shuffle.    *\/   private def prepareShuffleDependency(rdd: RDD[InternalRow],                                        outputAttributes: Seq[Attribute],                                        partitioning: OrderBucketsPartitioning,                                        serializer: Serializer,                                        part: Partitioner                                       )   : ShuffleDependency[OrderAndSortKey, InternalRow, InternalRow] = {     val rddWithKeys: RDD[Product2[OrderAndSortKey, InternalRow]] = {       rdd.mapPartitionsWithIndexInternal((i, iter) => {         val mutablePair = new MutablePair[OrderAndSortKey, InternalRow]()         val keyGen = getPartitionKeyExtractor(outputAttributes, partitioning, i)         iter.map { row => mutablePair.update(keyGen(row), row.copy) }       })     }       def keyOrdering[A &lt;: Product3[InternalRow, InternalRow, InternalRow]]: Ordering[A] = {       val allExprs = outputAttributes         .filter(attr => partitioning.sortExpressions.map(_.child.references).contains(attr.references))       implicit val sortOrdering: Ordering[InternalRow] =         new LazilyGeneratedOrdering(partitioning.sortExpressions, allExprs)       Ordering.by(_._3)     }       implicit val order: Ordering[OrderAndSortKey] = keyOrdering       \/\/ Now, we manually create a ShuffleDependency.     new ShuffleDependency[OrderAndSortKey, InternalRow, InternalRow](       rddWithKeys,       part,       serializer,       Some(order)     )   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0414\u043b\u044f \u0442\u043e\u0433\u043e, \u0447\u0442\u043e\u0431\u044b \u043d\u0435 \u043c\u0435\u043d\u044f\u0442\u044c \u0441\u0445\u0435\u043c\u0443 \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430 \u0438 \u0441\u043a\u0440\u044b\u0442\u044c \u043f\u043e\u0434\u0440\u043e\u0431\u043d\u043e\u0441\u0442\u0438 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0438 \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0438, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u0442\u0441\u044f \u043a\u043b\u0430\u0441\u0441\u00a0<em>org.apache.spark.sql.execution.exchange.ShuffledOrderRDD<\/em><\/p>\n<details class=\"spoiler\">\n<summary>ShuffledOrderRDD<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package org.apache.spark.sql.execution.exchange   import org.apache.spark.rdd.RDD import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.sql.execution.{CoalescedPartitioner, ShuffledRowRDDPartition} import org.apache.spark._ import ru.kalininskii.orderbucketing.OrderAndSortKey   \/**  * \u0421\u043f\u0435\u0446\u0438\u0430\u043b\u0438\u0437\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441-\u043d\u0430\u0441\u043b\u0435\u0434\u043d\u0438\u043a [[org.apache.spark.rdd.ShuffledRDD]]  * \u041f\u043e\u0434\u0434\u0435\u0440\u0436\u0438\u0432\u0430\u0435\u0442 \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u0443\u044e \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0443 \u043f\u043e \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u043c \u043f\u043e\u043b\u044f\u043c  *\/ class ShuffledOrderRDD(var dependency: ShuffleDependency[OrderAndSortKey, InternalRow, InternalRow])   extends RDD[InternalRow](dependency.rdd.context, Nil) {     private[this] val numPreShufflePartitions = dependency.partitioner.numPartitions     private[this] val partitionStartIndices: Array[Int] = (0 until numPreShufflePartitions).toArray     private[this] val part: Partitioner =     new CoalescedPartitioner(dependency.partitioner, partitionStartIndices)     override def getDependencies: Seq[Dependency[_]] = List(dependency)     override val partitioner: Option[Partitioner] = Some(part)     override def getPartitions: Array[Partition] = {     assert(partitionStartIndices.length == part.numPartitions)     Array.tabulate[Partition](partitionStartIndices.length) { i =>       val startIndex = partitionStartIndices(i)       val endIndex =         if (i &lt; partitionStartIndices.length - 1) {           partitionStartIndices(i + 1)         } else {           numPreShufflePartitions         }       new ShuffledRowRDDPartition(i, startIndex, endIndex)     }   }     override def getPreferredLocations(partition: Partition): Seq[String] = {     val tracker = SparkEnv.get.mapOutputTracker.asInstanceOf[MapOutputTrackerMaster]     val dep = dependencies.head.asInstanceOf[ShuffleDependency[_, _, _]]     tracker.getPreferredLocationsForShuffle(dep, partition.index)   }     override def compute(split: Partition, context: TaskContext): Iterator[InternalRow] = {     val shuffledRowPartition = split.asInstanceOf[ShuffledRowRDDPartition]     \/\/ The range of pre-shuffle partitions that we are fetching at here is     \/\/ [startPreShufflePartitionIndex, endPreShufflePartitionIndex - 1].     val reader =     SparkEnv.get.shuffleManager.getReader(       dependency.shuffleHandle,       shuffledRowPartition.startPreShufflePartitionIndex,       shuffledRowPartition.endPreShufflePartitionIndex,       context)     reader.read().asInstanceOf[Iterator[Product2[OrderAndSortKey, InternalRow]]].map(_._2)   }     override def clearDependencies() {     super.clearDependencies()     dependency = null   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0418 \u0432\u043e\u0442, \u0434\u043e\u043b\u0433\u043e\u0436\u0434\u0430\u043d\u043d\u044b\u0439 \u043a\u043b\u0430\u0441\u0441, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0431\u0443\u0434\u0435\u0442 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0442\u044c \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0438 \u043e\u0442\u044b\u0441\u043a\u0438\u0432\u0430\u0442\u044c \u043d\u043e\u043c\u0435\u0440 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD.<\/p>\n<p>\u041e\u0441\u043d\u043e\u0432\u043d\u044b\u043c \u043f\u043e\u043b\u0435\u043c \u0432 \u043d\u0451\u043c \u044f\u0432\u043b\u044f\u0435\u0442\u0441\u044f <strong>private<\/strong> <strong>var<\/strong> <em>rangeBounds<\/em>: Map[P, (Array[O], Int)].<br \/> \u042d\u0442\u043e \u043a\u0430\u0440\u0442\u0430, \u043a\u043b\u044e\u0447\u043e\u043c \u043a\u043e\u0442\u043e\u0440\u043e\u0439 \u044f\u0432\u043b\u044f\u0435\u0442\u0441\u044f InternalRow \u0441\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c\u0438 \u043f\u043e\u043b\u0435\u0439 \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f, \u0430 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435\u043c \u2013 \u043a\u043e\u0440\u0442\u0435\u0436 \u0438\u0437 \u043c\u0430\u0441\u0441\u0438\u0432\u0430 \u0432\u0435\u0440\u0445\u043d\u0438\u0445 \u0433\u0440\u0430\u043d\u0438\u0446 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439\u00a0RDD (\u0442\u043e\u0436\u0435 InternalRow, \u043d\u043e \u0441 \u0434\u0440\u0443\u0433\u0438\u043c \u0441\u043e\u0434\u0435\u0440\u0436\u0438\u043c\u044b\u043c) \u0438 \u0441\u0442\u0430\u0440\u0442\u043e\u0432\u043e\u0433\u043e \u043d\u043e\u043c\u0435\u0440\u0430 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 \u0434\u043b\u044f \u0441\u0435\u043a\u0446\u0438\u0438 (\u0432 \u0434\u0430\u043b\u044c\u043d\u0435\u0439\u0448\u0435\u043c \u043c\u043e\u0436\u0435\u0442 \u043d\u0430\u0437\u044b\u0432\u0430\u0442\u044c\u0441\u044f &#171;\u0441\u043c\u0435\u0449\u0435\u043d\u0438\u0435&#187;).<\/p>\n<p>\u0421\u044d\u043c\u043f\u043b\u0438\u043d\u0433 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u0441\u044f \u0441\u0442\u0430\u043d\u0434\u0430\u0440\u0442\u043d\u044b\u043c \u0441\u043f\u043e\u0441\u043e\u0431\u043e\u043c (RDD.sample), \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043c\u043e\u0436\u0435\u0442 \u0443\u043f\u0443\u0441\u043a\u0430\u0442\u044c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u0438\u0437 \u043c\u0430\u043b\u0435\u043d\u044c\u043a\u0438\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439. \u041d\u0430\u043f\u043e\u043c\u043d\u044e, \u0447\u0442\u043e \u043a\u043b\u0430\u0441\u0441 \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u043b\u0441\u044f \u0434\u043b\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0431\u043e\u043b\u044c\u0448\u0438\u0445 \u043e\u0431\u044a\u0435\u043c\u043e\u0432 \u0434\u0430\u043d\u043d\u044b\u0445, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043c\u043e\u0436\u0435\u0442 \u043d\u0435\u044d\u0444\u0444\u0435\u043a\u0442\u0438\u0432\u043d\u043e \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u043d\u0430 \u043d\u0435\u0441\u043a\u043e\u043b\u044c\u043a\u0438\u0445 \u0434\u0435\u0441\u044f\u0442\u043a\u0430\u0445 \u0438\u043b\u0438 \u0441\u043e\u0442\u043d\u044f\u0445 \u0441\u0442\u0440\u043e\u043a.<\/p>\n<p>\u041e\u0442\u0431\u043e\u0440 \u0433\u0440\u0430\u043d\u0438\u0446 \u0442\u0430\u043a\u0436\u0435 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u0441\u044f \u0432 RDD, \u043f\u0440\u0438 \u044d\u0442\u043e\u043c \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439 RDD \u0434\u043b\u044f \u043a\u0430\u0436\u0434\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438, \u044f\u0432\u043d\u043e\u0439 \u0438\u043b\u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0439,\u00a0\u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0435\u0442\u0441\u044f \u0434\u0438\u043d\u0430\u043c\u0438\u0447\u0435\u0441\u043a\u0438 \u043d\u0430 \u043e\u0441\u043d\u043e\u0432\u0430\u043d\u0438\u0438 \u0441\u043e\u0431\u0440\u0430\u043d\u043d\u043e\u0439 \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u0438 \u043e \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u0435 \u0441\u0442\u0440\u043e\u043a \u0432 \u0441\u044d\u043c\u043f\u043b\u0435.<\/p>\n<p>\u041f\u043e\u0441\u043b\u0435 \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u0438\u044f \u0432\u0441\u0435\u0445 \u043e\u0442\u043e\u0431\u0440\u0430\u043d\u043d\u044b\u0445 \u0433\u0440\u0430\u043d\u0438\u0446, \u043d\u0430 \u0434\u0440\u0430\u0439\u0432\u0435\u0440\u0435\u00a0\u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442\u0441\u044f \u0430\u043b\u0433\u043e\u0440\u0438\u0442\u043c \u0441\u043e \u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0438\u043c\u0438 \u0438\u043d\u0432\u0430\u0440\u0438\u0430\u043d\u0442\u0430\u043c\u0438:<\/p>\n<ol>\n<li>\n<p>\u041d\u0443\u043b\u0435\u0432\u0430\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f \u0440\u0435\u0437\u0435\u0440\u0432\u0438\u0440\u0443\u0435\u0442\u0441\u044f \u0434\u043b\u044f \u043d\u0435 \u0432\u043e\u0448\u0435\u0434\u0448\u0438\u0445 \u0432 \u0441\u0442\u0430\u0442\u0438\u0441\u0442\u0438\u043a\u0443 \u043a\u043b\u044e\u0447\u0435\u0439 (\u0441\u0435\u043a\u0446\u0438\u0439 Hive \u0438 \u043d\u0435\u044f\u0432\u043d\u044b\u0445 \u0441\u0435\u043a\u0446\u0438\u0439). \u0412 \u0441\u043b\u0443\u0447\u0430\u044f\u0445, \u043a\u043e\u0433\u0434\u0430 \u043d\u0443\u043b\u0435\u0432\u0430\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u0442 \u0441\u043b\u0438\u0448\u043a\u043e\u043c \u0431\u043e\u043b\u044c\u0448\u043e\u0439 \u043e\u0431\u044a\u0451\u043c \u0434\u0430\u043d\u043d\u044b\u0445, \u0438\u0437\u043c\u0435\u043d\u0438\u0442\u0435 \u043f\u043e\u043b\u044f \u044f\u0432\u043d\u043e\u0433\u043e \u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0433\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0438\u043b\u0438 \u043e\u0442\u043a\u0430\u0436\u0438\u0442\u0435\u0441\u044c \u043e\u0442 \u043d\u0438\u0445 \u043f\u043e\u043b\u043d\u043e\u0441\u0442\u044c\u044e;<\/p>\n<\/li>\n<li>\n<p>\u0417\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u044f\u0432\u043d\u044b\u0445 \u0438 \u043d\u0435\u044f\u0432\u043d\u044b\u0445 \u0441\u0435\u043a\u0446\u0438\u0439 \u043c\u043e\u0433\u0443\u0442 \u0441\u043b\u0435\u0434\u043e\u0432\u0430\u0442\u044c \u0432 \u043d\u0435\u043e\u0442\u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u043c \u043f\u043e\u0440\u044f\u0434\u043a\u0435, \u043d\u043e \u044d\u0442\u043e\u0442 \u043f\u043e\u0440\u044f\u0434\u043e\u043a \u0434\u043e\u043b\u0436\u0435\u043d \u0431\u044b\u043b \u0437\u0430\u0444\u0438\u043a\u0441\u0438\u0440\u043e\u0432\u0430\u043d \u0441 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435\u043c \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 &#171;\u0441\u043c\u0435\u0449\u0435\u043d\u0438\u044f&#187;. \u0421\u043c\u0435\u0449\u0435\u043d\u0438\u0435 &#8212; \u044d\u0442\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u0442\u0438\u043f\u0430 Int, \u043e\u043d\u043e \u0432\u0445\u043e\u0434\u0438\u0442 \u0432 \u043a\u043e\u0440\u0442\u0435\u0436Map[P, (Array[O], <em>Int<\/em>)];<\/p>\n<\/li>\n<li>\n<p>\u0421\u0442\u0430\u0440\u0442\u043e\u0432\u044b\u0439 \u043d\u043e\u043c\u0435\u0440 \u0441\u0435\u043a\u0446\u0438\u0438 \u2013 \u0441\u043c\u0435\u0449\u0435\u043d\u0438\u0435, \u043d\u0430\u0447\u0438\u043d\u0430\u0435\u0442\u0441\u044f \u0441 \u043e\u0434\u043d\u043e\u0433\u043e (\u0441\u043c. \u043f.1 \u2013 \u043d\u0443\u043b\u0435\u0432\u0430\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f \u0437\u0430\u0440\u0435\u0437\u0435\u0440\u0432\u0438\u0440\u043e\u0432\u0430\u043d\u0430, \u0438 \u0434\u043b\u044f \u043a\u0430\u0436\u0434\u043e\u0433\u043e \u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0435\u0433\u043e \u043a\u043b\u044e\u0447\u0430 \u043f\u0440\u0438\u0440\u0430\u0441\u0442\u0430\u0435\u0442 \u043d\u0430 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0433\u0440\u0430\u043d\u0438\u0446 + 1, \u0447\u0442\u043e \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0443\u0435\u0442 \u0440\u0435\u0430\u043b\u044c\u043d\u043e\u043c\u0443 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u0443 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439 \u0434\u043b\u044f \u0441\u0435\u043a\u0446\u0438\u0438 (\u0441\u043c. \u043f.4);<\/p>\n<\/li>\n<li>\n<p>\u041c\u0430\u0441\u0441\u0438\u0432 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f \u043d\u0435 \u0441\u043e\u0434\u0435\u0440\u0436\u0438\u0442 \u0441\u0430\u043c\u043e\u0439 \u043f\u043e\u0441\u043b\u0435\u0434\u043d\u0435\u0439 \u0432\u0435\u0440\u0445\u043d\u0435\u0439 \u0433\u0440\u0430\u043d\u0438\u0446\u044b, \u044d\u0442\u043e \u043b\u043e\u0433\u0438\u0447\u043d\u043e, \u0447\u0442\u043e\u0431\u044b \u043d\u0435 \u0441\u043e\u0437\u0434\u0430\u0432\u0430\u0442\u044c \u0432\u0441\u0435\u0433\u0434\u0430 \u043f\u0443\u0441\u0442\u0443\u044e \u0434\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u0443\u044e \u0441\u0435\u043a\u0446\u0438\u044e.<\/p>\n<\/li>\n<\/ol>\n<p>\u041f\u043e\u0438\u0441\u043a \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 \u043e\u0441\u0443\u0449\u0435\u0441\u0442\u0432\u043b\u044f\u0435\u0442\u0441\u044f \u0441\u043d\u0430\u0447\u0430\u043b\u0430 \u043f\u043e \u043a\u043b\u044e\u0447\u0443 \u043a\u0430\u0440\u0442\u044b (\u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0435 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u0432\u0441\u0435\u0445 \u043f\u043e\u043b\u0435\u0439 \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f), \u0437\u0430\u0442\u0435\u043c \u0432 \u043c\u0430\u0441\u0441\u0438\u0432\u0435 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439, \u043f\u043e\u0441\u043b\u0435\u0434\u043e\u0432\u0430\u0442\u0435\u043b\u044c\u043d\u044b\u043c \u0438\u043b\u0438 \u0434\u0432\u043e\u0438\u0447\u043d\u044b\u043c \u043f\u043e\u0438\u0441\u043a\u043e\u043c, \u0432 \u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e\u0441\u0442\u0438 \u043e\u0442 \u0434\u043b\u0438\u043d\u044b \u043c\u0430\u0441\u0441\u0438\u0432\u0430. \u0412\u0441\u0435 \u0437\u0430\u043f\u0438\u0441\u0438, \u0434\u043b\u044f \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u043a\u043b\u044e\u0447 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u0435\u0442 \u0432 \u043a\u0430\u0440\u0442\u0435, \u043e\u0442\u043f\u0440\u0430\u0432\u043b\u044f\u044e\u0442\u0441\u044f \u0432 \u043d\u0443\u043b\u0435\u0432\u0443\u044e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044e RDD, \u0447\u0442\u043e \u043c\u043e\u0436\u0435\u0442 \u043f\u0440\u0438\u0432\u0435\u0441\u0442\u0438 \u043a \u043d\u0435\u043f\u0440\u0438\u0435\u043c\u043b\u0435\u043c\u043e \u0431\u043e\u043b\u044c\u0448\u043e\u043c\u0443 \u043e\u0431\u044a\u0451\u043c\u0443 \u044d\u0442\u043e\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438, \u0435\u0441\u043b\u0438 \u043f\u043e\u043b\u044f \u044f\u0432\u043d\u043e\u0433\u043e \u0438\u043b\u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0433\u043e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0432\u044b\u0431\u0440\u0430\u043d\u044b \u043d\u0435\u0443\u0434\u0430\u0447\u043d\u043e \u0438 \u0438\u043c\u0435\u044e\u0442 \u0441\u043b\u0438\u0448\u043a\u043e\u043c \u0432\u044b\u0441\u043e\u043a\u0443\u044e \u0441\u0435\u043b\u0435\u043a\u0442\u0438\u0432\u043d\u043e\u0441\u0442\u044c. \u0412 \u044d\u0442\u043e\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u043b\u0443\u0447\u0448\u0435 \u043d\u0435 \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u0442\u044c \u0442\u0430\u0431\u043b\u0438\u0446\u0443 \u0432\u043e\u043e\u0431\u0449\u0435, \u043e\u0441\u0442\u0430\u0432\u0438\u0432 \u0442\u043e\u043b\u044c\u043a\u043e \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435 \u043d\u0430 \u0431\u0430\u043a\u0435\u0442\u044b \u043f\u043e \u043f\u043e\u043b\u044e \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f.<\/p>\n<p>\u0412\u044b\u0431\u043e\u0440 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u043e\u0434\u043d\u043e\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD, \u043d\u0430\u0447\u0438\u043d\u0430\u044f \u0441 \u043f\u0440\u043e\u0442\u043e\u0442\u0438\u043f\u0430, \u043e\u0441\u0443\u0449\u0435\u0441\u0442\u0432\u043b\u044f\u0435\u0442\u0441\u044f \u0442\u0430\u043a: \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043c\u0435\u043d\u044c\u0448\u0435 \u043b\u0438\u0431\u043e \u0440\u0430\u0432\u043d\u044b\u0435 \u0432\u0435\u0440\u0445\u043d\u0435\u0439 \u0433\u0440\u0430\u043d\u0438\u0446\u0435 (\u043e\u043d\u0430, \u043a\u0430\u043a \u043c\u044b \u043f\u043e\u043c\u043d\u0438\u043c, \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0430 \u0437\u0430\u0440\u0430\u043d\u0435\u0435 \u0438 \u0438\u0437\u0432\u0435\u0441\u0442\u043d\u0430), \u0438 \u0431\u043e\u043b\u044c\u0448\u0435, \u0447\u0435\u043c \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0430\u044f \u0432\u0435\u0440\u0445\u043d\u044f\u044f \u0433\u0440\u0430\u043d\u0438\u0446\u0430, \u0438\u043b\u0438 \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0430\u044f \u0432\u0435\u0440\u0445\u043d\u044f\u044f \u0433\u0440\u0430\u043d\u0438\u0446\u0430 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u0443\u0435\u0442 (\u044d\u0442\u043e \u043f\u0435\u0440\u0432\u044b\u0439 \u0444\u0430\u0439\u043b \u0432 \u044f\u0432\u043d\u043e\u0439 \u0438\u043b\u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438).<\/p>\n<p>\u042f \u0437\u043d\u0430\u043a\u043e\u043c \u0441 \u043a\u043e\u043d\u0446\u0435\u043f\u0446\u0438\u0435\u0439 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u044c\u043d\u043e\u0433\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0432 Oracle, \u0438 \u043c\u0435\u0436\u0434\u0443 \u044d\u0442\u0438\u043c\u0438 \u0434\u0432\u0443\u043c\u044f \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u044f\u043c\u0438 \u0435\u0441\u0442\u044c \u043e\u0442\u043b\u0438\u0447\u0438\u0435 \u0432 \u0441\u0442\u0440\u043e\u0433\u043e\u0441\u0442\u0438 \u043d\u0435\u0440\u0430\u0432\u0435\u043d\u0441\u0442\u0432. \u0412 Oracle \u0441\u0435\u043a\u0446\u0438\u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u044e\u0442\u0441\u044f \u043f\u043e \u0443\u0441\u043b\u043e\u0432\u0438\u044e \u00ab\u043c\u0435\u043d\u044c\u0448\u0435 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u00bb \u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u043c\u0443 \u0443\u0441\u043b\u043e\u0432\u0438\u044e \u00ab\u0431\u043e\u043b\u044c\u0448\u0435 \u0438\u043b\u0438 \u0440\u0430\u0432\u043d\u043e \u043f\u0440\u0435\u0434\u044b\u0434\u0443\u0449\u0435\u0433\u043e \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u00bb.<\/p>\n<p>\u042d\u0442\u043e \u0441\u0432\u044f\u0437\u0430\u043d\u043e \u0441 \u0442\u0435\u043c, \u0447\u0442\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 Oracle \u0440\u0430\u0431\u043e\u0442\u0430\u0435\u0442 \u0441 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u044b\u043c\u0438 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c\u0438: \u0434\u043b\u044f \u0432\u0440\u0435\u043c\u0435\u043d\u0438 \u044d\u0442\u043e \u0447\u0430\u0449\u0435 \u0432\u0441\u0435\u0433\u043e \u0441\u0443\u0442\u043a\u0438, \u0434\u043b\u044f \u0447\u0438\u0441\u043b\u043e\u0432\u044b\u0445 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u2013 \u043c\u0438\u043b\u043b\u0438\u043e\u043d\u044b \u0438\u043b\u0438 \u043c\u0438\u043b\u043b\u0438\u0430\u0440\u0434\u044b \u0438 \u0442\u0430\u043a \u0434\u0430\u043b\u0435\u0435.<\/p>\n<p>\u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u043e\u0447\u0435\u043d\u044c \u0443\u0434\u043e\u0431\u043d\u043e, \u043a \u043f\u0440\u0438\u043c\u0435\u0440\u0443, \u0447\u0442\u043e \u043c\u043e\u0436\u043d\u043e \u0443\u043a\u0430\u0437\u0430\u0442\u044c \u0434\u043b\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u044f \u043f\u043e \u0441\u0443\u0442\u043a\u0430\u043c:\u00a0values less than \u20182021-12-01 00:00:00\u2019 \u2013 \u043f\u0440\u0438 \u044d\u0442\u043e\u043c \u0432\u0438\u0434\u043d\u043e, \u0447\u0442\u043e \u0432 \u043d\u043e\u0432\u0443\u044e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044e \u043f\u043e\u043f\u0430\u0434\u0451\u0442 \u043d\u0430\u0447\u0430\u043b\u043e \u0441\u0443\u0442\u043e\u043a (\u0438 \u043c\u0435\u0441\u044f\u0446\u0430), \u0430 \u0441\u043a\u043e\u043b\u044c \u0443\u0433\u043e\u0434\u043d\u043e \u0431\u043b\u0438\u0437\u043a\u043e\u0435, \u043d\u043e \u043c\u0435\u043d\u044c\u0448\u0435\u0435 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u0431\u0443\u0434\u0435\u0442 \u0432 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438, \u043e\u0442\u043d\u043e\u0441\u044f\u0449\u0435\u0439\u0441\u044f \u043a \u20182021-11-30\u2019.<\/p>\n<p>\u0412 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u0434\u0435\u043b\u043e \u0441\u043e\u0432\u0441\u0435\u043c \u0432 \u0434\u0440\u0443\u0433\u043e\u043c: \u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u0435\u043b\u044c \u043d\u0438\u0447\u0435\u0433\u043e \u043d\u0435 \u0434\u043e\u043b\u0436\u0435\u043d \u0437\u043d\u0430\u0442\u044c \u043e\u0431 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u0445 \u0438 \u0444\u0430\u0439\u043b\u0430\u0445, \u0435\u0433\u043e \u0438\u043d\u0442\u0435\u0440\u0435\u0441\u0443\u044e\u0442 \u0442\u043e\u043b\u044c\u043a\u043e \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u044b\u0435 \u0434\u0430\u043d\u043d\u044b\u0435 \u0438 \u044d\u0442\u043e \u0437\u0430\u0434\u0430\u0447\u0430 \u044f\u0432\u043d\u043e\u0433\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f (\u043f\u0435\u0440\u0432\u044b\u0439 \u0443\u0440\u043e\u0432\u0435\u043d\u044c).<\/p>\n<p>\u0422\u0430\u043a\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c, \u0437\u0430\u0434\u0430\u0447\u0430 \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u044f \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u2013 \u043e\u0431\u0435\u0441\u043f\u0435\u0447\u0438\u0442\u044c \u043f\u0440\u0438\u043c\u0435\u0440\u043d\u043e \u0440\u0430\u0432\u043d\u044b\u0439 \u0440\u0430\u0437\u043c\u0435\u0440 \u0444\u0430\u0439\u043b\u043e\u0432, \u043b\u0451\u0433\u043a\u043e\u0441\u0442\u044c \u043d\u0430\u0445\u043e\u0436\u0434\u0435\u043d\u0438\u044f \u043d\u0443\u0436\u043d\u044b\u0445 \u0444\u0430\u0439\u043b\u043e\u0432 \u0438 \u0441\u043e\u0432\u0435\u0440\u0448\u0435\u043d\u043d\u043e \u043f\u0440\u043e\u0437\u0440\u0430\u0447\u043d\u043e\u0435 \u0432\u0437\u0430\u0438\u043c\u043e\u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0435, \u0435\u0441\u043b\u0438 \u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u0435\u043b\u044c \u043e\u0431\u0440\u0430\u0449\u0430\u0435\u0442\u0441\u044f \u043d\u0430\u043f\u0440\u044f\u043c\u0443\u044e \u043a \u0442\u0440\u0430\u0434\u0438\u0446\u0438\u043e\u043d\u043d\u044b\u043c \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0430\u043c. \u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u044e\u0442\u0441\u044f \u00ab\u043d\u0430 \u043b\u0435\u0442\u0443\u00bb, \u0445\u0440\u0430\u043d\u044f\u0442\u0441\u044f \u0432 \u043e\u0442\u0434\u0435\u043b\u044c\u043d\u043e\u0439 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0435 \u0438 \u0441\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u044b \u0441 \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u0438\u043c\u0438 \u0444\u0430\u0439\u043b\u0430\u043c\u0438. \u041a\u0430\u0436\u0434\u0430\u044f \u0441\u0435\u043a\u0446\u0438\u044f \u0438\u043c\u0435\u0435\u0442 \u0441\u0432\u043e\u0439 \u0438\u043d\u0434\u0438\u0432\u0438\u0434\u0443\u0430\u043b\u044c\u043d\u044b\u0439 \u043d\u0430\u0431\u043e\u0440 \u0433\u0440\u0430\u043d\u0438\u0446. \u041a\u0440\u043e\u043c\u0435 \u0442\u043e\u0433\u043e, \u0443\u043a\u0430\u0436\u0435\u043c \u043d\u0430 \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0432\u0430\u0436\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442: \u0435\u0441\u043b\u0438 \u043e\u0434\u043d\u0443 \u0438\u0437 \u0433\u0440\u0430\u043d\u0438\u0446, \u0432 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435, \u0432\u0435\u0440\u0445\u043d\u044e\u044e, \u043c\u044b \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u044f\u0435\u043c \u044f\u0432\u043d\u043e, \u0442\u043e \u0434\u0440\u0443\u0433\u0430\u044f, \u043d\u0438\u0436\u043d\u044f\u044f \u0433\u0440\u0430\u043d\u0438\u0446\u0430 \u043d\u0435 \u043e\u0447\u0435\u043d\u044c \u0432\u0430\u0436\u043d\u0430, \u043d\u043e \u043b\u0443\u0447\u0448\u0435, \u0447\u0442\u043e\u0431\u044b \u0435\u0451 \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u043b\u0435\u0433\u043a\u043e \u043d\u0430\u0439\u0442\u0438, \u0435\u0441\u043b\u0438 \u0432 \u0434\u0430\u043b\u044c\u043d\u0435\u0439\u0448\u0435\u043c \u043e\u043d\u0430 \u0431\u0443\u0434\u0435\u0442 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u043e\u0432\u0430\u0442\u044c \u0432 \u043a\u0430\u0440\u0442\u0435 \u0434\u0430\u043d\u043d\u044b\u0445. \u0414\u043b\u044f \u043d\u0438\u0436\u043d\u0435\u0439 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u044d\u0442\u043e \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u0442\u0430\u043a, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0435 \u043d\u0443\u0436\u043d\u043e \u043f\u0440\u043e\u0447\u0438\u0442\u0430\u0442\u044c \u0438\u0437 \u043d\u0430\u0447\u0430\u043b\u0430 \u0444\u0430\u0439\u043b\u0430 \u0438\u043b\u0438 \u0438\u0437 \u043d\u0430\u0447\u0430\u043b\u0430 \u0441\u0435\u0433\u043c\u0435\u043d\u0442\u0430 \u0434\u0430\u043d\u043d\u044b\u0445, \u0432 \u0441\u043b\u0443\u0447\u0430\u0435 \u043a\u043e\u043b\u043e\u043d\u043e\u0447\u043d\u043e\u0433\u043e \u0444\u043e\u0440\u043c\u0430\u0442\u0430 \u0445\u0440\u0430\u043d\u0435\u043d\u0438\u044f.<\/p>\n<details class=\"spoiler\">\n<summary>OrderBucketsPartitioner<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package org.apache.spark.sql.partitioning   import java.io.{IOException, ObjectInputStream, ObjectOutputStream}   import org.apache.spark.{Partitioner, SparkEnv} import org.apache.spark.rdd.RDD import org.apache.spark.serializer.JavaSerializer import org.apache.spark.sql.catalyst.InternalRow import org.apache.spark.util.{CollectionsUtils, Utils}   import scala.collection.mutable.ArrayBuffer import scala.reflect.ClassTag import scala.util.hashing.byteswap32   \/**  * \u041d\u0430\u0441\u043b\u0435\u0434\u043d\u0438\u043a [[org.apache.spark.Partitioner]] \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u0435\u0442 \u043f\u043e \u0442\u043e\u0447\u043d\u044b\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043d\u0435\u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u0439 (\u043f\u043e\u043b\u0435\u0439).  * \u0418 \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u043e\u0432 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043e\u0434\u043d\u043e\u0433\u043e \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f (\u043f\u043e\u043b\u044f)  *  * @note The actual number of partitions created by the RangePartitioner might not be the same  *       as the `partitions` parameter, in the case where the number of sampled records is less than  *       the value of `partitions`.  *\/ class OrderBucketsPartitioner[P: ClassTag, O: Ordering : ClassTag ](rdd: RDD[_ &lt;: Product2[P, O]],   val numLines: Int,   val partitions: Int,   val distribution: Option[Seq[(Int, P, Either[(O, Seq[String]), (O, O, Int)])]])   extends Partitioner {     \/\/ We allow partitions = 0, which happens when sorting an empty RDD under the default settings.   require(partitions >= 0, s\"Number of partitions cannot be negative but found $partitions.\")     private var ordering = implicitly[Ordering[O]]     private def getDistribution: Option[Map[P, (Array[O], Int)]] = distribution flatMap {     case Seq() => None     case distr =>       val (inner, outer) = distr.partition(_._3.isLeft)         \/\/\u0432\u043d\u0435\u0448\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0438\u0437 \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f.       val (singleFileBounds, manyFileBounds): (Seq[(P, (Int, (O, O, Int)))], Seq[(P, (Int, (O, O, Int)))]) =         outer           .map { case (index, p, right) => (p, (index, right.right.get)) }           .partition { case (_, (_, (_, _, filesNum))) => filesNum == 1 }       \/\/\u041d\u0443\u0436\u043d\u043e \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u0432\u043d\u0443\u0442\u0440\u0435\u043d\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b, \u0442\u043e\u043b\u044c\u043a\u043e \u0434\u043b\u044f \u0442\u0435\u0445, \u0433\u0434\u0435 \u0431\u043e\u043b\u044c\u0448\u0435 \u043e\u0434\u043d\u043e\u0433\u043e \u0444\u0430\u0439\u043b\u0430       val outerBounds: Map[P, Seq[(Int, (O, O, Int))]] = manyFileBounds         .groupBy(_._1)         .map { case (p, values) => p -> values.map(bounds => bounds._2) }       \/\/ \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0435\u0434\u0438\u043d\u0438\u0447\u043d\u044b\u0445 \u0444\u0430\u0439\u043b\u043e\u0432 \u0443\u0436\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u044b       val singleBounds: Array[(P, Array[(Int, Array[O])])] = singleFileBounds         .groupBy(_._1)         .map { case (p, values) =>           (p, values.map(bounds => (bounds._2._1, Array(bounds._2._2._2))).toArray) }         .toArray         \/\/\u0432\u043d\u0443\u0442\u0440\u0435\u043d\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u043f\u043e\u043b\u0443\u0447\u0430\u044e\u0442\u0441\u044f \u0438\u0437 rdd \u0444\u0438\u043b\u044c\u0442\u0440\u043e\u043c, \u0441\u044d\u043c\u043f\u043b\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435\u043c, \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u043e\u0439       val innerBounds: Array[(P, Array[(Int, Array[O])])] =         OrderBucketsPartitioner.getInnerBounds(rdd, numLines, partitions, outerBounds)         \/\/\u0437\u0430\u0440\u0430\u043d\u0435\u0435 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0451\u043d\u043d\u044b\u0435 \u0432\u043d\u0443\u0442\u0440\u0435\u043d\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b       val givenBounds: Array[(P, Array[(Int, Array[O])])] = inner         .map { case (index, p, left) => (p, (index, Array(left.left.get._1))) }         .groupBy(_._1)         .map { case (p, values) => p -> values.map(bounds => bounds._2).toArray }         .toArray         \/\/\u0442\u0435\u043f\u0435\u0440\u044c \u043e\u0431\u0430 \u043c\u0430\u0441\u0441\u0438\u0432\u0430 \u0441\u043a\u043b\u0435\u0438\u0432\u0430\u044e\u0442\u0441\u044f, \u0433\u0440\u0443\u043f\u043f\u0438\u0440\u0443\u044e\u0442\u0441\u044f, \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c \u0432\u0441\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b,       \/\/ \u043e\u0442\u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u044b\u0435 \u043f\u043e \u0432\u043e\u0437\u0440\u0430\u0441\u0442\u0430\u043d\u0438\u044e, \u0431\u0435\u0437 \u0441\u0430\u043c\u043e\u0439 \u0432\u0435\u0440\u0445\u043d\u0435\u0439,       val fixedBounds: Map[P, (Array[O], Int)] = (givenBounds ++ innerBounds ++ singleBounds)         .groupBy(_._1)         .map { case (p, values) =>           val (indexes, bounds) = values.flatMap(_._2).unzip           p -> (bounds.flatten.sorted.init,             indexes.min)         }         Some(fixedBounds)   }     \/\/ \u041a\u0430\u0440\u0442\u0430 \u0441 \u043c\u0430\u0441\u0441\u0438\u0432\u0430\u043c\u0438 \u0432\u0435\u0440\u0445\u043d\u0438\u0445 \u0433\u0440\u0430\u043d\u0438\u0446 \u0434\u043b\u044f \u043f\u0435\u0440\u0432\u044b\u0445 (partitions - 1) \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439   private var rangeBounds: Map[P, (Array[O], Int)] = getDistribution.getOrElse {     if (partitions &lt;= 1) {       Map.empty[P, (Array[O], Int)]     } else {       val groupedBounds = OrderBucketsPartitioner.getBounds(rdd, numLines, partitions)         \/\/\u0438\u043d\u0434\u0435\u043a\u0441\u044b \u0441\u043c\u0435\u0449\u0435\u043d\u0438\u0439, \u0447\u0442\u043e\u0431\u044b \u043e\u0431\u0435\u0441\u043f\u0435\u0447\u0438\u0442\u044c \u0443\u043d\u0438\u043a\u0430\u043b\u044c\u043d\u044b\u0439 \u043d\u043e\u043c\u0435\u0440 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438, \u0435\u0441\u043b\u0438 \u043c\u0430\u0441\u0441\u0438\u0432 \u0433\u0440\u0430\u043d\u0438\u0446 \u043f\u0443\u0441\u0442\u043e\u0439, \u0442\u043e \u044d\u0442\u043e \u043e\u0434\u043d\u0430 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f       val overallIndex = groupedBounds.scanLeft(1) { case (sum, (_, orders)) => sum + math.max(orders.length + 1, 1) }         groupedBounds         .zip(overallIndex)         .map { case ((part, orders), index) => (part, (orders, index)) }         .toMap     }   }     def getRangeBounds: Map[P, (Array[O], Int)] = rangeBounds     \/\/\u0432\u0435\u0441\u044c \u043c\u0430\u0441\u0441\u0438\u0432 \u043f\u043b\u044e\u0441 \u043e\u0434\u043d\u0430 \u043d\u0443\u043b\u0435\u0432\u0430\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f \u0434\u043b\u044f \u043d\u0435 \u043f\u043e\u043f\u0430\u0432\u0448\u0438\u0445 \u0432 \u0441\u044d\u043c\u043f\u043b \u0442\u0430\u0441\u043a\u043e\u0432   \/\/ \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438, \u043d\u0435 \u043f\u043e\u043f\u0430\u0432\u0448\u0438\u0435 \u0432 \u0441\u044d\u043c\u043f\u043b\u044b \u0434\u043e\u043b\u0436\u043d\u044b \u0431\u044b\u0442\u044c \u043d\u0435\u0431\u043e\u043b\u044c\u0448\u0438\u043c\u0438, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043d\u0438\u0447\u0435\u0433\u043e \u0441\u0442\u0440\u0430\u0448\u043d\u043e\u0433\u043e \u043d\u0435 \u0431\u0443\u0434\u0435\u0442   \/\/ \u0434\u0430\u0436\u0435 \u0435\u0441\u043b\u0438 \u043e\u043d\u0438 \u043f\u043e\u043f\u0430\u0434\u0443\u0442 \u0432 \u043e\u0434\u0438\u043d \u0442\u0430\u0441\u043a, \u0442\u0430\u043a \u0431\u0443\u0434\u0435\u0442 \u0434\u0430\u0436\u0435 \u043b\u0443\u0447\u0448\u0435 - \u043c\u0435\u043d\u044c\u0448\u0435 \u0444\u0430\u0439\u043b\u043e\u0432   def numPartitions: Int = rangeBounds.values.map(_._1.length + 1).sum + 1     private var binarySearch: (Array[O], O) => Int = CollectionsUtils.makeBinarySearch[O]     def getPartition(key: Any): Int = {     val (p, k, _) = key.asInstanceOf[(P, O, InternalRow)]     val (bounds, shift) = rangeBounds.getOrElse(p, (Array.empty[O], 0))     var partition = 0     if (bounds.length &lt; 16) {       \/\/ If we have less than 16 partitions naive search       while (partition &lt; bounds.length &amp;&amp; ordering.gt(k, bounds(partition))) {         partition += 1       }     } else {       \/\/ \u043c\u0435\u0442\u043e\u0434 \u0431\u0438\u043d\u0430\u0440\u043d\u043e\u0433\u043e \u043f\u043e\u0438\u0441\u043a\u0430 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0451\u043d \u0442\u043e\u043b\u044c\u043a\u043e \u043e\u0434\u0438\u043d \u0440\u0430\u0437       partition = binarySearch(bounds, k)       \/\/ binarySearch either returns the match location or -[insertion point]-1       if (partition &lt; 0) {         partition = -partition - 1       }       if (partition > bounds.length) {         partition = bounds.length       }     }     shift + partition   }     override def equals(other: Any): Boolean = other match {     case r: OrderBucketsPartitioner[_, _] =>       r.rangeBounds == rangeBounds     case _ =>       false   }     override def hashCode(): Int = {     val prime = 31     var result = 1     var i = 0     val arr = rangeBounds.values.map(_._1).toArray.flatten     while (i &lt; arr.length) {       result = prime * result + arr(i).hashCode()       i += 1     }     result = prime * result     result   }     @throws(classOf[IOException])   private def writeObject(out: ObjectOutputStream): Unit = Utils.tryOrIOException {     val sfactory = SparkEnv.get.serializer     sfactory match {       case _: JavaSerializer => out.defaultWriteObject()       case _ =>         out.writeObject(ordering)         out.writeObject(binarySearch)           val ser = sfactory.newInstance()         Utils.serializeViaNestedStream(out, ser) { stream =>           stream.writeObject(scala.reflect.classTag[Map[P, (Array[O], Int)]])           stream.writeObject(rangeBounds)         }     }   }     @throws(classOf[IOException])   private def readObject(in: ObjectInputStream): Unit = Utils.tryOrIOException {     val sfactory = SparkEnv.get.serializer     sfactory match {       case _: JavaSerializer => in.defaultReadObject()       case _ =>         ordering = in.readObject().asInstanceOf[Ordering[O]]         binarySearch = in.readObject().asInstanceOf[(Array[O], O) => Int]           val ser = sfactory.newInstance()         Utils.deserializeViaNestedStream(in, ser) { ds =>           implicit val classTag: ClassTag[Map[P, (Array[O], Int)]] = ds.readObject[ClassTag[Map[P, (Array[O], Int)]]]()           rangeBounds = ds.readObject[Map[P, (Array[O], Int)]]()         }     }   } }   object OrderBucketsPartitioner {     \/**    * \u0420\u0430\u0441\u0441\u0447\u0438\u0442\u044b\u0432\u0430\u0435\u0442 \u043a\u043e\u044d\u0444\u0444\u0438\u0446\u0438\u0435\u043d\u0442 \u0441\u044d\u043c\u043b\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f \u043e\u0431\u0440\u0430\u0442\u043d\u0443\u044e \u043b\u043e\u0433\u0430\u0440\u0438\u0444\u043c\u0438\u0447\u0435\u0441\u043a\u0443\u044e \u043f\u0440\u043e\u043f\u043e\u0440\u0446\u0438\u044e,    * \u0447\u0435\u043c \u0431\u043e\u043b\u044c\u0448\u0435 \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445, \u0442\u0435\u043c \u043c\u0435\u043d\u044c\u0448\u0435 \u0447\u0430\u0441\u0442\u044c, \u043a\u043e\u0442\u043e\u0440\u0430\u044f \u0431\u0443\u0434\u0435\u0442 \u0432\u0437\u044f\u0442\u0430 \u0438\u0437 \u043d\u0435\u0433\u043e    *    * @param numLines - \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u0444\u0430\u0439\u043b\u0435 - \u043e\u043f\u0430\u0441\u043d\u043e\u0435 \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u043e\u0436\u0435\u043d\u0438\u0435, \u043d\u043e \u0441\u0434\u0435\u043b\u0430\u0442\u044c \u0435\u0433\u043e \u043c\u043e\u0436\u043d\u043e    * @param numParts - \u043f\u0440\u0438\u0431\u043b\u0438\u0437\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u043e\u0446\u0435\u043d\u0435\u043d\u043d\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439    *\/   private def getSampleCoefficient(numLines: Int, numParts: Int): Double = {     val numLinesD = numLines.toDouble     (1e4 \/ (numLinesD * math.log(numLinesD * numParts))) min 1.0   }     \/**    * \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u044f\u0435\u0442 \u0441\u044d\u043c\u043f\u043b RDD    *    * @param rdd               - \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445    * @param sampleCoefficient - 0 &lt; \u043a\u043e\u044d\u0444\u0444\u0438\u0446\u0438\u0435\u043d\u0442 &lt;= 1.0    *\/   private def getSampleRDD[O: Ordering : ClassTag,     P: ClassTag](rdd: RDD[_ &lt;: Product2[P, O]],                  sampleCoefficient: Double): RDD[_ &lt;: Product2[P, O]] = {     val sampledRdd: RDD[_ &lt;: Product2[P, O]] = if (sampleCoefficient &lt; 1.0) {       val seed = byteswap32(-rdd.id - 1)       rdd.sample(withReplacement = false, sampleCoefficient, seed)     } else {       rdd     }     sampledRdd   }     \/**    * \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442\u0441\u044f \u0432 \u043a\u0430\u0436\u0434\u043e\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 \u0441\u044d\u043c\u043b\u043f\u0430 RDD    * \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u0442 \u0432\u0435\u0440\u0445\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0434\u043b\u044f \u043a\u0430\u0436\u0434\u043e\u0433\u043e \u0444\u0430\u0439\u043b\u0430 (\u0442\u0430\u0441\u043a\u0430, \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD)    *    * @param iter     \u043d\u0435\u043e\u0442\u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u044b\u0439 \u0438\u0442\u0435\u0440\u0430\u0442\u043e\u0440 \u0441\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c\u0438 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439 \u0438 \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f    * @param initStep \u0448\u0430\u0433, \u0440\u0430\u0432\u043d\u044b\u0439 \u0436\u0435\u043b\u0430\u0435\u043c\u043e\u043c\u0443 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u0443 \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435    * @param weight   \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0445 \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u043d\u0430 \u0441\u0442\u0440\u043e\u043a\u0443 \u0441\u044d\u043c\u043f\u043b\u0430    * @return \u0433\u0440\u0430\u043d\u0438\u0446\u044b    *\/   private def determineBoundsWithinPartition[P: ClassTag, O: Ordering : ClassTag   ](iter: Iterator[(P, O)],     weight: Double,     initStep: Int): Iterator[(P, Array[O])] = {     prepareMap(iter)       .map { case (p, ordered) =>         val partitions: Int = math.ceil(ordered.length * weight \/ initStep).intValue()         val bounds: Array[O] = extractBounds(ordered, partitions, weight)         (p, bounds)       }.iterator   }     private def prepareMap[O: Ordering : ClassTag, P: ClassTag](iter: Iterator[(P, O)]): Map[P, Array[O]] = {     val mappedParts = iter.toArray       .groupBy(_._1)       .map { case (p, opArray) => (p, opArray.map(_._2).sorted) }     mappedParts   }     \/\/\u043e\u0442\u0431\u043e\u0440 \u0441\u0430\u043c\u0438\u0445 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439   private def extractBounds[O: Ordering : ClassTag](values: Array[O],                                                     partitions: Int,                                                     weight: Double)(implicit ordering: Ordering[O]): Array[O] = {     val numCandidates = values.length     val partCount = numCandidates * weight     var cumWeight = 0.0     val step: Int = (partCount \/ partitions).intValue()     var target = step     val innerBounds = ArrayBuffer.empty[O]     var i = 0     var j = 0     var prevInnerBound = Option.empty[O]       while ((i &lt; numCandidates) &amp;&amp; (j &lt; partitions - 1)) {       val key = values(i)       cumWeight += weight       if (cumWeight >= target) {         \/\/ Skip duplicate values.         if (prevInnerBound.isEmpty || ordering.gt(key, prevInnerBound.get)) {           innerBounds += key           target += step           j += 1           prevInnerBound = Some(key)         }       }       i += 1     }     innerBounds.toArray   }     \/**    * \u0441\u044d\u043c\u043f\u043b\u0438\u0440\u0443\u0435\u0442 RDD \u0438 \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u0442 \u043a\u0430\u0440\u0442\u0443, \u0433\u0434\u0435 \u043a\u043b\u044e\u0447\u0438 - \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043f\u043e\u043b\u0435\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f,    * \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f - \u043c\u0430\u0441\u0441\u0438\u0432\u044b \u0441 \u0432\u0435\u0440\u0445\u043d\u0438\u043c\u0438 \u0433\u0440\u0430\u043d\u0438\u0446\u0430\u043c\u0438 \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f    *    * @param rdd      \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445    * @param numLines \u0436\u0435\u043b\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435    * @param numParts \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439    * @return \u0433\u0440\u0430\u043d\u0438\u0446\u044b    *\/   private def getBounds[P: ClassTag, O: Ordering : ClassTag](rdd: RDD[_ &lt;: Product2[P, O]],                                                              numLines: Int,                                                              numParts: Int): Array[(P, Array[O])] = {     val sampleCoefficient = getSampleCoefficient(numLines, numParts) \/\/rdd.partitions.length       val sampledRdd: RDD[_ &lt;: Product2[P, O]] = getSampleRDD(rdd, sampleCoefficient)       val weight = 1.0 \/ sampleCoefficient     sampledRdd       .map(row => (row._1, row._2))       .partitionBy(new Partitioner {         override def numPartitions: Int = numParts           override def getPartition(key: Any): Int = Utils.nonNegativeMod(key.hashCode, numParts)       })       .mapPartitions(iter => determineBoundsWithinPartition(iter, weight, numLines))       .collect()   }       \/**    * \u0444\u0438\u043b\u044c\u0442\u0440\u0443\u0435\u0442 RDD \u043d\u0430 \u0440\u0430\u0432\u0435\u043d\u0441\u0442\u0432\u043e \u043f\u043e\u043b\u0435\u0439 \u043a \u043d\u0443\u0436\u043d\u044b\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u0438 \u043f\u043e \u043d\u0430\u0445\u043e\u0436\u0434\u0435\u043d\u0438\u044e \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u0432\u043d\u0435\u0448\u043d\u0438\u0445 \u0433\u0440\u0430\u043d\u0438\u0446    *    * @param rdd    \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445    * @param bounds \u0433\u0440\u0430\u043d\u0438\u0446\u044b RDD    * @return \u0442\u0449\u0430\u0442\u0435\u043b\u044c\u043d\u043e \u043e\u0442\u0444\u0438\u043b\u044c\u0442\u0440\u043e\u0432\u0430\u043d\u043d\u044b\u0439 RDD    *\/   private def filterRDDByOuterBounds[P: ClassTag, O: Ordering : ClassTag   ](rdd: RDD[_ &lt;: Product2[P, O]],     bounds: Map[P, Seq[(Int, (O, O, Int))]]): RDD[_ &lt;: Product2[P, O]] = {     val ordering = implicitly[Ordering[O]]     \/\/\u043f\u043e\u043b\u0443\u0447\u0430\u0435\u043c Option[Array[Tuple3]] \u0438 \u043f\u0440\u043e\u0432\u0435\u0440\u044f\u0435\u043c exists     rdd       .filter { item =>         bounds           .get(item._1)           .exists(_.exists(b => ordering.lteq(b._2._1, item._2) &amp;&amp; ordering.lteq(item._2, b._2._2)))       }   }     \/**    * \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442\u0441\u044f \u0432 \u043a\u0430\u0436\u0434\u043e\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 \u0441\u044d\u043c\u043b\u043f\u0430 RDD    * \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u0442 \u0432\u0435\u0440\u0445\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0434\u043b\u044f \u043a\u0430\u0436\u0434\u043e\u0433\u043e \u0444\u0430\u0439\u043b\u0430 (\u0442\u0430\u0441\u043a\u0430, \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438 RDD)    *    * @param iter     \u043d\u0435\u043e\u0442\u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u044b\u0439 \u0438\u0442\u0435\u0440\u0430\u0442\u043e\u0440 \u0441\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c\u0438 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439 \u0438 \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f    * @param initStep \u0448\u0430\u0433, \u0440\u0430\u0432\u043d\u044b\u0439 \u0436\u0435\u043b\u0430\u0435\u043c\u043e\u043c\u0443 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u0443 \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435    * @param weight   \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0440\u0435\u0430\u043b\u044c\u043d\u044b\u0445 \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u043d\u0430 \u0441\u0442\u0440\u043e\u043a\u0443 \u0441\u044d\u043c\u043f\u043b\u0430    * @return \u0433\u0440\u0430\u043d\u0438\u0446\u044b    *\/   private def determineInnerBoundsWithinPartition[P: ClassTag, O: Ordering : ClassTag   ](iter: Iterator[(P, O)],     weight: Double,     initStep: Int,     outerBounds: Map[P, Seq[(Int, (O, O, Int))]]): Iterator[(P, Array[(Int, Array[O])])] = {     val ordering = implicitly[Ordering[O]]       prepareMap(iter)       .map { case (p, orderedValues) =>         (p, outerBounds(p).toArray.map { case (index, (lowerO, upperO, partitions)) =>           val ordered = orderedValues             .dropWhile(ordering.lt(_, lowerO))             .takeWhile(ordering.lteq(_, upperO))             val innerBounds: Array[O] = extractBounds(ordered, partitions, weight)             (index, innerBounds :+ upperO)         })       }.iterator   }     \/**    * \u0441\u044d\u043c\u043f\u043b\u0438\u0440\u0443\u0435\u0442 RDD \u0438 \u043f\u043e\u043b\u0443\u0447\u0430\u0435\u0442 \u043a\u0430\u0440\u0442\u0443, \u0433\u0434\u0435 \u043a\u043b\u044e\u0447\u0438 - \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043f\u043e\u043b\u0435\u0439 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f,    * \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f - \u043c\u0430\u0441\u0441\u0438\u0432\u044b \u0441 \u0432\u0435\u0440\u0445\u043d\u0438\u043c\u0438 \u0433\u0440\u0430\u043d\u0438\u0446\u0430\u043c\u0438 \u043f\u043e\u043b\u044f \u0443\u043f\u043e\u0440\u044f\u0434\u043e\u0447\u0438\u0432\u0430\u043d\u0438\u044f    * \u0412\u0430\u0436\u043d\u043e, \u0447\u0442\u043e\u0431\u044b \u043c\u043e\u0436\u043d\u043e \u0431\u044b\u043b\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u0432\u0435\u0440\u0445\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0438,    * \u0447\u0442\u043e\u0431\u044b \u043d\u0435 \u0434\u043e\u043f\u0443\u0441\u0442\u0438\u0442\u044c \u043f\u043e\u0433\u043b\u043e\u0449\u0435\u043d\u0438\u044f \u0432\u043e\u0437\u043c\u043e\u0436\u043d\u044b\u0445 \u0432\u043d\u0443\u0442\u0440\u0435\u043d\u043d\u0438\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439    *    * @param rdd      \u043e\u0441\u043d\u043e\u0432\u043d\u043e\u0439 \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445    * @param numLines \u0436\u0435\u043b\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435    * @param numParts \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439    * @param bounds   \u0432\u043d\u0435\u0448\u043d\u0438\u0435 \u0433\u0440\u0430\u043d\u0438\u0446\u044b RDD    * @return \u0433\u0440\u0430\u043d\u0438\u0446\u044b    *\/   private def getInnerBounds[P: ClassTag, O: Ordering : ClassTag](rdd: RDD[_ &lt;: Product2[P, O]],                                                                   numLines: Int,                                                                   numParts: Int,                                                                   bounds: Map[P, Seq[(Int, (O, O, Int))]]                                                                  ): Array[(P, Array[(Int, Array[O])])] = {     val sampleCoefficient = getSampleCoefficient(numLines, numParts) \/\/rdd.partitions.length       val filterRdd: RDD[_ &lt;: Product2[P, O]] = filterRDDByOuterBounds(rdd, bounds)     val sampledRdd: RDD[_ &lt;: Product2[P, O]] = getSampleRDD(filterRdd, sampleCoefficient)       val weight = 1.0 \/ sampleCoefficient     sampledRdd       .map(row => (row._1, row._2))       .partitionBy(new Partitioner {         override def numPartitions: Int = numParts           override def getPartition(key: Any): Int = Utils.nonNegativeMod(key.hashCode, bounds.size)       })       .mapPartitions(iter => determineInnerBoundsWithinPartition(iter, weight, numLines, bounds))       .collect()   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0427\u0442\u043e\u0431\u044b \u0441\u0432\u044f\u0437\u0430\u0442\u044c \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u0438\u0439 \u043f\u043b\u0430\u043d \u0438 \u0441\u0430\u043c\u043e \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u0435, \u043d\u0430\u043c \u043d\u0443\u0436\u043d\u043e \u0431\u0443\u0434\u0435\u0442 \u0441\u043e\u0437\u0434\u0430\u0442\u044c \u043e\u0431\u044a\u0435\u043a\u0442, \u0440\u0430\u0441\u0448\u0438\u0440\u044f\u044e\u0449\u0438\u0439 SparkStrategy.<\/p>\n<details class=\"spoiler\">\n<summary>RepartitionStrategy<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.rangebucketing.strategy   import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan import org.apache.spark.sql.execution.{SparkPlan, SparkStrategy, exchange} import ru.kalininskii.rangebucketing.plans.logical.RepartitionByRangeBuckets   object RepartitionStrategy extends SparkStrategy {   override def apply(plan: LogicalPlan): Seq[SparkPlan] = plan match {     case r: RepartitionByRangeBuckets =>       exchange.ShuffleExchangeBucketsExec(r.partitioning, planLater(r.child)) :: Nil     case _ => Nil   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0435\u043f\u0435\u0440\u044c \u043c\u043e\u0436\u043d\u043e \u0441\u0434\u0435\u043b\u0430\u0442\u044c \u0438\u043d\u044a\u0435\u043a\u0446\u0438\u044e \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d\u0438\u044f \u0432 SparkSession:<\/p>\n<details class=\"spoiler\">\n<summary>Extension injection<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>spark = SparkSession   .builder()   .appName(\"PARTITION_TEST\")   .master(\"local[*]\")   .config(sparkConf)   .withExtensions(e => {     e.injectPlannerStrategy(_ => RepartitionStrategy)   })   .getOrCreate()<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0415\u0441\u043b\u0438 \u0438\u043d\u044a\u0435\u043a\u0446\u0438\u0438 \u0438\u043b\u0438 \u0441\u043f\u043e\u0441\u043e\u0431\u0430 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440\u0430 \u043d\u0435 \u0431\u0443\u0434\u0435\u0442, \u0442\u043e \u043c\u044b \u0443\u0432\u0438\u0434\u0438\u043c \u0442\u0430\u043a\u0443\u044e \u043e\u0448\u0438\u0431\u043a\u0443:<\/p>\n<details class=\"spoiler\">\n<summary>\u0421\u043e\u043e\u0431\u0449\u0435\u043d\u0438\u0435 \u043e\u0431 \u043e\u0448\u0438\u0431\u043a\u0435<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>assertion failed: No plan for RepartitionWithOrderAndSort [ts_part#55], id#13 ASC NULLS FIRST, 4000, 24 +- Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, substring(cast(ts#18 as string), 0, 10) AS ts_part#55]    +- InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18], StorageLevel(disk, memory, deserialized, 1 replicas)          +- LocalTableScan [id#13, name#14, amount#15, value#16, divider#17, ts#18]  java.lang.AssertionError: assertion failed: No plan for RepartitionWithOrderAndSort [ts_part#55], id#13 ASC NULLS FIRST, 4000, 24 +- Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, substring(cast(ts#18 as string), 0, 10) AS ts_part#55]    +- InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18], StorageLevel(disk, memory, deserialized, 1 replicas)          +- LocalTableScan [id#13, name#14, amount#15, value#16, divider#17, ts#18]           at scala.Predef$.assert(Predef.scala:170)          at org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:93)          at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:78)          at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2$$anonfun$apply$2.apply(QueryPlanner.scala:75)          at scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157)          at scala.collection.TraversableOnce$$anonfun$foldLeft$1.apply(TraversableOnce.scala:157)          at scala.collection.Iterator$class.foreach(Iterator.scala:891)          at scala.collection.AbstractIterator.foreach(Iterator.scala:1334)          at scala.collection.TraversableOnce$class.foldLeft(TraversableOnce.scala:157)          at scala.collection.AbstractIterator.foldLeft(Iterator.scala:1334)          at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:75)          at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$2.apply(QueryPlanner.scala:67)          at scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:435)          at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:441)          at org.apache.spark.sql.catalyst.planning.QueryPlanner.plan(QueryPlanner.scala:93)          at org.apache.spark.sql.execution.QueryExecution.sparkPlan$lzycompute(QueryExecution.scala:72)          at org.apache.spark.sql.execution.QueryExecution.sparkPlan(QueryExecution.scala:68)          at org.apache.spark.sql.execution.QueryExecution.executedPlan$lzycompute(QueryExecution.scala:77)          at org.apache.spark.sql.execution.QueryExecution.executedPlan(QueryExecution.scala:77)          at org.apache.spark.sql.execution.QueryExecution$$anonfun$toString$3.apply(QueryExecution.scala:207)          at org.apache.spark.sql.execution.QueryExecution$$anonfun$toString$3.apply(QueryExecution.scala:207)          at org.apache.spark.sql.execution.QueryExecution.stringOrError(QueryExecution.scala:99)          at org.apache.spark.sql.execution.QueryExecution.toString(QueryExecution.scala:207)          at org.apache.spark.sql.execution.command.ExplainCommand.run(commands.scala:167)          at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)          at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)          at org.apache.spark.sql.execution.command.ExecutedCommandExec.executeCollect(commands.scala:79)          at org.apache.spark.sql.Dataset.explain(Dataset.scala:485)          at ru.kalininskii.rangebucketing.RangeBucketingPartitionTest$$anonfun$1.apply(RangeBucketingPartitionTest.scala:49)          at ru.kalininskii.rangebucketing.RangeBucketingPartitionTest$$anonfun$1.apply(RangeBucketingPartitionTest.scala:42)<\/code><\/pre>\n<\/div>\n<\/details>\n<p><em>\u0421\u0433\u0435\u043d\u0435\u0440\u0438\u0440\u0443\u0435\u043c \u043d\u0430\u0431\u043e\u0440 \u0434\u0430\u043d\u043d\u044b\u0445 \u0441 \u043f\u0441\u0435\u0432\u0434\u043e\u0441\u043b\u0443\u0447\u0430\u0439\u043d\u044b\u043c\u0438 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c\u0438 \u0438 \u043f\u0440\u0438\u043c\u0435\u043d\u0438\u043c \u043a \u043d\u0435\u043c\u0443 \u0441\u043e\u0437\u0434\u0430\u043d\u043d\u044b\u0439 \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440<\/em>:\u00a0<\/p>\n<details class=\"spoiler\">\n<summary>\u0414\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>    val dataFrameSource = getRandomDF(0, rangeLimit, 3, 0).persist()       .withColumn(\"ts_part\", fn.expr(\"substring(ts,0,10)\"))       dataFrameSource.printSchema()       val dfRep = dataFrameSource       .repartitionWithOrderAndSort(numLines, rangeLimit \/ numLines,         fn.col(\"event_time\"), List(fn.col(\"ts_part\")), List(fn.col(\"id\")))       .persist()       dfRep.explain(true)<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0412 \u043f\u043b\u0430\u043d\u0430\u0445 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f \u0432\u0438\u0434\u043d\u044b \u043d\u043e\u0432\u044b\u0435 \u043f\u0443\u043d\u043a\u0442\u044b, \u043d\u0430 \u0432\u0441\u0435\u0445 \u044d\u0442\u0430\u043f\u0430\u0445 \u0438\u0445 \u043c\u043e\u0436\u043d\u043e \u043e\u0442\u0441\u043b\u0435\u0434\u0438\u0442\u044c<\/p>\n<details class=\"spoiler\">\n<summary>\u041f\u043b\u0430\u043d \u0437\u0430\u043f\u0440\u043e\u0441\u0430 \u0441 \u043f\u0440\u0438\u043c\u0435\u043d\u0435\u043d\u0438\u0435\u043c \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440\u0430<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>== Parsed Logical Plan == 'Repartition with unknown distribution: ['ts_part], 'id ASC NULLS FIRST, 4000, 24 +- 'Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, 'substring('ts, 0, 10) AS ts_part#55]    +- Project [_1#6 AS id#13, _2#7 AS name#14, _3#8 AS amount#15, _4#9 AS value#16, _5#10 AS divider#17, _6#11 AS ts#18]       +- LocalRelation [_1#6, _2#7, _3#8, _4#9, _5#10, _6#11]  == Analyzed Logical Plan == id: int, name: string, amount: decimal(38,18), value: decimal(38,18), divider: string, ts: timestamp, ts_part: string Repartition with unknown distribution: [ts_part#55], id#13 ASC NULLS FIRST, 4000, 24 +- Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, substring(cast(ts#18 as string), 0, 10) AS ts_part#55]    +- Project [_1#6 AS id#13, _2#7 AS name#14, _3#8 AS amount#15, _4#9 AS value#16, _5#10 AS divider#17, _6#11 AS ts#18]       +- LocalRelation [_1#6, _2#7, _3#8, _4#9, _5#10, _6#11]  == Optimized Logical Plan == InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18, ts_part#55], StorageLevel(disk, memory, deserialized, 1 replicas)    +- ExchangeByRangeBuckets repartition with unknown distribution:(ts_part#55, id#13 ASC NULLS FIRST, 4000, 24, None)       +- *(1) Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, substring(cast(ts#18 as string), 0, 10) AS ts_part#55]          +- InMemoryTableScan [amount#15, divider#17, id#13, name#14, ts#18, value#16]                +- InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18], StorageLevel(disk, memory, deserialized, 1 replicas)                      +- LocalTableScan [id#13, name#14, amount#15, value#16, divider#17, ts#18]  == Physical Plan == InMemoryTableScan [id#13, name#14, amount#15, value#16, divider#17, ts#18, ts_part#55]    +- InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18, ts_part#55], StorageLevel(disk, memory, deserialized, 1 replicas)          +- ExchangeByRangeBuckets repartition with unknown distribution:(ts_part#55, id#13 ASC NULLS FIRST, 4000, 24, None)             +- *(1) Project [id#13, name#14, amount#15, value#16, divider#17, ts#18, substring(cast(ts#18 as string), 0, 10) AS ts_part#55]                +- InMemoryTableScan [amount#15, divider#17, id#13, name#14, ts#18, value#16]                      +- InMemoryRelation [id#13, name#14, amount#15, value#16, divider#17, ts#18], StorageLevel(disk, memory, deserialized, 1 replicas)                            +- LocalTableScan [id#13, name#14, amount#15, value#16, divider#17, ts#18]<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0435\u0441\u0442\u044b \u0434\u043e\u043b\u0436\u043d\u044b \u043f\u043e\u043a\u0430\u0437\u0430\u0442\u044c, \u0447\u0442\u043e \u043a\u043b\u0430\u0441\u0441\u044b \u043f\u043b\u0430\u043d\u043e\u0432 \u043d\u0430 \u0441\u0430\u043c\u043e\u043c \u0434\u0435\u043b\u0435 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044e\u0442\u0441\u044f (\u043f\u0440\u043e\u0441\u0442\u044b\u0435 \u044e\u043d\u0438\u0442-\u0442\u0435\u0441\u0442\u044b \u043f\u043b\u0430\u043d\u043e\u0432 Spark):<\/p>\n<details class=\"spoiler\">\n<summary>\u0418\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0435 \u043a\u043b\u0430\u0441\u0441\u043e\u0432 <\/summary>\n<div class=\"spoiler__content\">\n<pre><code>    val logicalPlanString = \"Repartition plan with unknown distribution: \" +       \"partition by ['ts_part], \" +       \"global order by 'event_time, \" +       \"local sort by ['id], \" +       \"4000 rows per task, 24 initial tasks\"       val analyzedPlanString = \"Repartition plan with unknown distribution: \" +       \"partition by [ts_part#64], \" +       \"global order by event_time#21, \" +       \"local sort by [id#15], \" +       \"4000 rows per task, 24 initial tasks\"       val sparkPlanString = \"ExchangeWithOrder Repartition with unknown distribution: \" +       \"partition by [ts_part#64], \" +       \"global order by event_time#21, \" +       \"local sort by [id#15], \" +       \"4000 rows per task, 24 initial tasks\"       assert(dfRep.queryExecution.logical.toString contains logicalPlanString)     assert(dfRep.queryExecution.analyzed.toString contains analyzedPlanString)     assert(dfRep.queryExecution.sparkPlan.toString contains sparkPlanString)<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439 \u0432\u0441\u0435\u0433\u0434\u0430 \u0431\u0443\u0434\u0435\u0442 \u0440\u0430\u0432\u043d\u043e 28<\/p>\n<details class=\"spoiler\">\n<summary>\u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439 RDD<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>println(dfRep.rdd.getNumPartitions) 28 dfRep.rdd.getNumPartitions should be(28)<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0412\u044b\u0432\u0435\u0434\u0435\u043c \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c \u0441 \u043d\u043e\u043c\u0435\u0440\u0430\u043c\u0438 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439, \u043d\u0443\u043b\u0435\u0432\u0430\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f \u0434\u043e\u043b\u0436\u043d\u0430 \u043e\u0442\u0441\u0443\u0442\u0441\u0442\u0432\u043e\u0432\u0430\u0442\u044c, \u043f\u043e\u0442\u043e\u043c\u0443 \u0447\u0442\u043e \u043d\u0435 \u0441\u043e\u0434\u0435\u0440\u0436\u0438\u0442 \u043d\u0438 \u043e\u0434\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f. \u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u2013 \u043a\u043e\u043b\u043e\u043d\u043a\u0430 \u00abcnt_\u00bb &#8212; \u0432\u0441\u0435\u0433\u0434\u0430 \u043d\u0435\u043c\u043d\u043e\u0433\u043e \u043c\u0435\u043d\u044c\u0448\u0435 \u0438 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0434\u043e\u0441\u0442\u0430\u0442\u043e\u0447\u043d\u043e \u0440\u0430\u0432\u043d\u043e\u043c\u0435\u0440\u043d\u043e\u0435, \u043e\u0442\u043a\u043b\u043e\u043d\u0435\u043d\u0438\u0435, \u043a\u0430\u043a \u043f\u0440\u0430\u0432\u0438\u043b\u043e, \u043e\u0441\u0442\u0430\u0451\u0442\u0441\u044f \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043f\u044f\u0442\u0438 \u043f\u0440\u043e\u0446\u0435\u043d\u0442\u043e\u0432.<\/p>\n<details class=\"spoiler\">\n<summary>\u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u044f\u0445 RDD<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>+----------+-----------+-----------------------+-----------------------+----+ |ts_part   |rdd_part_id|range_min_event_time   |range_max_event_time   |cnt_| +----------+-----------+-----------------------+-----------------------+----+ |2021-09-05|10         |2021-09-05 21:57:56.186|2021-09-05 21:57:56.454|3672| |2021-09-05|11         |2021-09-05 21:57:56.455|2021-09-05 21:57:56.547|3400| |2021-09-05|12         |2021-09-05 21:57:56.548|2021-09-05 21:57:56.642|3761| |2021-09-05|13         |2021-09-05 21:57:56.643|2021-09-06 21:57:56.447|3601| |2021-09-05|14         |2021-09-06 21:57:56.448|2021-09-06 21:57:56.544|3601| |2021-09-05|15         |2021-09-06 21:57:56.545|2021-09-06 21:57:56.641|3856| |2021-09-05|16         |2021-09-06 21:57:56.642|2021-09-07 21:57:56.449|3711| |2021-09-05|17         |2021-09-07 21:57:56.45 |2021-09-07 21:57:56.549|3703| |2021-09-05|18         |2021-09-07 21:57:56.55 |2021-09-07 21:57:56.647|3776| |2021-09-06|19         |2021-09-05 21:57:56.186|2021-09-05 21:57:56.454|3763| |2021-09-06|20         |2021-09-05 21:57:56.455|2021-09-05 21:57:56.564|3916| |2021-09-06|21         |2021-09-05 21:57:56.565|2021-09-06 21:57:56.241|3557| |2021-09-06|22         |2021-09-06 21:57:56.242|2021-09-06 21:57:56.459|3721| |2021-09-06|23         |2021-09-06 21:57:56.46 |2021-09-06 21:57:56.563|3714| |2021-09-06|24         |2021-09-06 21:57:56.564|2021-09-07 21:57:56.246|3520| |2021-09-06|25         |2021-09-07 21:57:56.247|2021-09-07 21:57:56.455|3671| |2021-09-06|26         |2021-09-07 21:57:56.456|2021-09-07 21:57:56.556|3716| |2021-09-06|27         |2021-09-07 21:57:56.557|2021-09-07 21:57:56.647|3670| |2021-09-07|1          |2021-09-05 21:57:56.188|2021-09-05 21:57:56.454|3690| |2021-09-07|2          |2021-09-05 21:57:56.455|2021-09-05 21:57:56.556|3704| |2021-09-07|3          |2021-09-05 21:57:56.557|2021-09-05 21:57:56.644|3601| |2021-09-07|4          |2021-09-05 21:57:56.645|2021-09-06 21:57:56.45 |3685| |2021-09-07|5          |2021-09-06 21:57:56.451|2021-09-06 21:57:56.55 |3767| |2021-09-07|6          |2021-09-06 21:57:56.551|2021-09-06 21:57:56.644|3846| |2021-09-07|7          |2021-09-06 21:57:56.645|2021-09-07 21:57:56.452|3782| |2021-09-07|8          |2021-09-07 21:57:56.453|2021-09-07 21:57:56.555|3775| |2021-09-07|9          |2021-09-07 21:57:56.556|2021-09-07 21:57:56.647|3820| +----------+-----------+-----------------------+-----------------------+----+<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0421\u043e\u0445\u0440\u0430\u043d\u0438\u043c \u0438 \u0432\u044b\u0432\u0435\u0434\u0435\u043c, \u0447\u0442\u043e\u0431\u044b \u0443\u0431\u0435\u0434\u0438\u0442\u044c\u0441\u044f, \u0447\u0442\u043e \u0432\u0441\u0435 \u0444\u0430\u0439\u043b\u044b \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0443\u044e\u0442 \u0438\u0441\u0445\u043e\u0434\u043d\u044b\u043c\u0438 \u0441\u0435\u043a\u0446\u0438\u044f\u043c<\/p>\n<details class=\"spoiler\">\n<summary>\u041f\u0440\u043e\u0432\u0435\u0440\u043a\u0430 \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0438\u044f<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>dfRep       .withColumn(\"rdd_part_id\", fn.spark_partition_id())       .groupBy(fn.col(\"ts_part\") as \"ts_part\", fn.col(\"rdd_part_id\"))       .agg(fn.min(\"event_time\") as \"range_min_event_time\",         fn.max(\"event_time\") as \"range_max_event_time\",         fn.count(fn.lit(0)) as \"cnt_\")       .orderBy(\"ts_part\", \"rdd_part_id\")       .show(100, false)       dfRep       .sortWithinPartitions(fn.col(\"id\"))       .write       .mode(SaveMode.Overwrite)       .partitionBy(\"ts_part\")       .format(\"parquet\")       .option(\"path\", \"\/test\/pa\/sometable\/snp\")       .saveAsTable(\"testschema.sometable\")       val dfT = spark.table(\"testschema.sometable\")       println(dfT.count)       dfT       .withColumn(\"file_name\",         fn.regexp_replace(fn.input_file_name(), \".*\/test\/pa\/sometable\/snp\", \"\"))       .groupBy(fn.col(\"ts_part\") as \"ts_part\", fn.col(\"file_name\"))       .agg(fn.min(\"event_time\") as \"range_min_event_time\",         fn.max(\"event_time\") as \"range_max_event_time\",         fn.count(fn.lit(0)) as \"cnt_\")       .orderBy(\"file_name\")       .show(100, false) <\/code><\/pre>\n<pre><code>+----------+---------------------------------------------------------------------------------------+-----------------------+-----------------------+----+ |ts_part   |file_name                                                                              |range_min_event_time   |range_max_event_time   |cnt_| +----------+---------------------------------------------------------------------------------------+-----------------------+-----------------------+----+ |2021-09-05|\/ts_part=2021-09-05\/part-00010-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.186|2021-09-05 21:57:56.454|3672| |2021-09-05|\/ts_part=2021-09-05\/part-00011-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.455|2021-09-05 21:57:56.547|3400| |2021-09-05|\/ts_part=2021-09-05\/part-00012-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.548|2021-09-05 21:57:56.642|3761| |2021-09-05|\/ts_part=2021-09-05\/part-00013-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.643|2021-09-06 21:57:56.447|3601| |2021-09-05|\/ts_part=2021-09-05\/part-00014-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.448|2021-09-06 21:57:56.544|3601| |2021-09-05|\/ts_part=2021-09-05\/part-00015-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.545|2021-09-06 21:57:56.641|3856| |2021-09-05|\/ts_part=2021-09-05\/part-00016-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.642|2021-09-07 21:57:56.449|3711| |2021-09-05|\/ts_part=2021-09-05\/part-00017-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.45 |2021-09-07 21:57:56.549|3703| |2021-09-05|\/ts_part=2021-09-05\/part-00018-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.55 |2021-09-07 21:57:56.647|3776| |2021-09-06|\/ts_part=2021-09-06\/part-00019-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.186|2021-09-05 21:57:56.454|3763| |2021-09-06|\/ts_part=2021-09-06\/part-00020-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.455|2021-09-05 21:57:56.564|3916| |2021-09-06|\/ts_part=2021-09-06\/part-00021-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.565|2021-09-06 21:57:56.241|3557| |2021-09-06|\/ts_part=2021-09-06\/part-00022-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.242|2021-09-06 21:57:56.459|3721| |2021-09-06|\/ts_part=2021-09-06\/part-00023-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.46 |2021-09-06 21:57:56.563|3714| |2021-09-06|\/ts_part=2021-09-06\/part-00024-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.564|2021-09-07 21:57:56.246|3520| |2021-09-06|\/ts_part=2021-09-06\/part-00025-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.247|2021-09-07 21:57:56.455|3671| |2021-09-06|\/ts_part=2021-09-06\/part-00026-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.456|2021-09-07 21:57:56.556|3716| |2021-09-06|\/ts_part=2021-09-06\/part-00027-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.557|2021-09-07 21:57:56.647|3670| |2021-09-07|\/ts_part=2021-09-07\/part-00001-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.188|2021-09-05 21:57:56.454|3690| |2021-09-07|\/ts_part=2021-09-07\/part-00002-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.455|2021-09-05 21:57:56.556|3704| |2021-09-07|\/ts_part=2021-09-07\/part-00003-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.557|2021-09-05 21:57:56.644|3601| |2021-09-07|\/ts_part=2021-09-07\/part-00004-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-05 21:57:56.645|2021-09-06 21:57:56.45 |3685| |2021-09-07|\/ts_part=2021-09-07\/part-00005-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.451|2021-09-06 21:57:56.55 |3767| |2021-09-07|\/ts_part=2021-09-07\/part-00006-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.551|2021-09-06 21:57:56.644|3846| |2021-09-07|\/ts_part=2021-09-07\/part-00007-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-06 21:57:56.645|2021-09-07 21:57:56.452|3782| |2021-09-07|\/ts_part=2021-09-07\/part-00008-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.453|2021-09-07 21:57:56.555|3775| |2021-09-07|\/ts_part=2021-09-07\/part-00009-28321a07-711d-4936-b305-680199f90649.c000.snappy.parquet|2021-09-07 21:57:56.556|2021-09-07 21:57:56.647|3820| +----------+---------------------------------------------------------------------------------------+-----------------------+-----------------------+----+<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u00a0\u041a\u0430\u043a \u0432\u0438\u0434\u0438\u043c, \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u044b \u0441\u043e\u0445\u0440\u0430\u043d\u0438\u043b\u0438\u0441\u044c, \u0434\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u043d\u0435 \u043f\u043e\u0442\u0440\u0435\u0431\u043e\u0432\u0430\u043b\u043e\u0441\u044c \u043d\u0438\u043a\u0430\u043a\u0438\u0445 \u0434\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u044b\u0445 \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0439. \u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0444\u0430\u0439\u043b\u043e\u0432 \u0434\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0441\u043b\u0443\u0447\u0430\u044f \u0432\u0441\u0435\u0433\u0434\u0430 \u043d\u0430 \u043e\u0434\u0438\u043d \u043c\u0435\u043d\u044c\u0448\u0435, \u0447\u0435\u043c \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439 \u0432 RDD. \u041e\u0434\u043d\u0430\u043a\u043e \u0432 \u043e\u0431\u0449\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435 \u044d\u0442\u043e \u043d\u0435 \u0442\u0430\u043a, \u0432\u0435\u0434\u044c \u0432 \u043d\u0443\u043b\u0435\u0432\u0443\u044e \u0441\u0435\u043a\u0446\u0438\u044e \u043c\u043e\u0433\u0443\u0442 \u043f\u043e\u043f\u0430\u0441\u0442\u044c \u0434\u0430\u043d\u043d\u044b\u0435 \u0438\u0437 \u043e\u0447\u0435\u043d\u044c \u043d\u0435\u0431\u043e\u043b\u044c\u0448\u0438\u0445 \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u0438\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439.\u00a0dfT.inputFiles.length should <em>be<\/em>(dfRep.<em>rdd<\/em>.getNumPartitions &#8212; 1)<\/p>\n<p>\u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u043e\u0441\u0442\u0430\u043b\u043e\u0441\u044c \u043f\u0440\u0435\u0436\u043d\u0438\u043c:<\/p>\n<details class=\"spoiler\">\n<summary>\u041a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u0441\u043e\u0445\u0440\u0430\u043d\u0435\u043d\u043d\u044b\u0445 \u0444\u0430\u0439\u043b\u0430\u0445<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>val savedCount = dfT.count savedCount shouldEqual rangeLimit<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0418 \u0441\u0445\u0435\u043c\u0430 \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430 \u043d\u0435 \u0438\u0437\u043c\u0435\u043d\u0438\u043b\u0430\u0441\u044c, \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u043d\u0435 \u043c\u0435\u043d\u044f\u0435\u0442 \u043f\u043e\u0440\u044f\u0434\u043e\u043a \u043f\u043e\u043b\u0435\u0439 \u0438 \u0438\u0445 \u0442\u0438\u043f\u044b.<\/p>\n<details class=\"spoiler\">\n<summary>\u0421\u0445\u0435\u043c\u044b \u0441\u043e\u0432\u043f\u0430\u0434\u0430\u044e\u0442<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>dfT.printSchema()<\/code><\/pre>\n<pre><code>root  |-- id: integer (nullable = true)  |-- name: string (nullable = true)  |-- amount: decimal(38,18) (nullable = true)  |-- value: decimal(38,18) (nullable = true)  |-- divider: string (nullable = true)  |-- ts: timestamp (nullable = true)  |-- event_time: timestamp (nullable = true)  |-- ts_part: string (nullable = true)<\/code><\/pre>\n<pre><code>dfT.schema should be (dataFrameSource.schema)<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u00a0\u0414\u043b\u044f \u0442\u043e\u0433\u043e, \u0447\u0442\u043e\u0431\u044b \u0443\u0431\u0435\u0434\u0438\u0442\u044c\u0441\u044f, \u0447\u0442\u043e \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0430 \u0440\u0430\u0431\u043e\u0442\u0430\u0435\u0442, \u0432\u044b\u0432\u0435\u0434\u0435\u043c \u0432\u0441\u0435 \u0437\u0430\u043f\u0438\u0441\u0438 \u0441 &#171;id&#187; \u043c\u0435\u043d\u044c\u0448\u0438\u043c\u00a032, \u0431\u0435\u0437 \u0434\u043e\u043f\u043e\u043b\u043d\u0438\u0442\u0435\u043b\u044c\u043d\u043e\u0439 \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0438. \u041c\u044b \u0443\u0432\u0438\u0434\u0438\u043c, \u0447\u0442\u043e \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u0444\u0430\u0439\u043b\u0430 \u0437\u0430\u043f\u0438\u0441\u0438 \u043e\u0442\u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u044b \u043f\u043e \u0432\u043e\u0437\u0440\u0430\u0441\u0442\u0430\u043d\u0438\u044e \u043f\u043e\u043b\u044f &#171;id&#187;, \u043d\u0435\u0441\u043c\u043e\u0442\u0440\u044f \u043d\u0430 \u0442\u043e, \u0447\u0442\u043e \u0431\u044b\u043b\u0438 \u0441\u043f\u0435\u0446\u0438\u0430\u043b\u044c\u043d\u043e \u043f\u0435\u0440\u0435\u043c\u0435\u0448\u0430\u043d\u044b \u043f\u0440\u0438 \u0441\u043e\u0437\u0434\u0430\u043d\u0438\u0438 \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430.<\/p>\n<details class=\"spoiler\">\n<summary>\u041f\u0440\u043e\u0432\u0435\u0440\u043a\u0430 \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u043e\u0439 \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0438<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>dfT   .withColumn(\"file_name\",     fn.regexp_replace(fn.input_file_name(), \".*\/test\/pa\/sometable\/snp\", \"\"))   .select(\"id\", \"file_name\")   .where(\"id &lt; 32\")   .show(32, false)<\/code><\/pre>\n<pre><code>+---+---------------------------------------------------------------------------------------+ |id |file_name                                                                              | +---+---------------------------------------------------------------------------------------+ |9  |\/ts_part=2021-09-06\/part-00022-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |15 |\/ts_part=2021-09-06\/part-00022-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |19 |\/ts_part=2021-09-06\/part-00022-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |25 |\/ts_part=2021-09-07\/part-00004-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |4  |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |13 |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |14 |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |16 |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |28 |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |31 |\/ts_part=2021-09-07\/part-00006-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |8  |\/ts_part=2021-09-05\/part-00012-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |12 |\/ts_part=2021-09-05\/part-00012-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |21 |\/ts_part=2021-09-05\/part-00012-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |30 |\/ts_part=2021-09-05\/part-00012-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |10 |\/ts_part=2021-09-05\/part-00015-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |27 |\/ts_part=2021-09-05\/part-00015-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |0  |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |1  |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |3  |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |20 |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |22 |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |26 |\/ts_part=2021-09-06\/part-00024-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |18 |\/ts_part=2021-09-07\/part-00001-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |5  |\/ts_part=2021-09-05\/part-00010-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |6  |\/ts_part=2021-09-05\/part-00010-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |23 |\/ts_part=2021-09-05\/part-00010-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |24 |\/ts_part=2021-09-05\/part-00010-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |29 |\/ts_part=2021-09-05\/part-00010-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |2  |\/ts_part=2021-09-06\/part-00019-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |7  |\/ts_part=2021-09-06\/part-00019-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |11 |\/ts_part=2021-09-06\/part-00019-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| |17 |\/ts_part=2021-09-06\/part-00019-ec8ffbb5-474a-4ce1-bcdb-cc5665800031.c000.snappy.parquet| +---+---------------------------------------------------------------------------------------+<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0414\u043e\u0431\u0430\u0432\u043b\u044e \u043f\u043e\u043b\u043d\u044b\u0435 \u0442\u0435\u0441\u0442\u043e\u0432\u044b\u0435 \u043a\u043b\u0430\u0441\u0441\u044b \u0438 pom.xml \u0434\u043b\u044f \u0441\u0431\u043e\u0440\u043a\u0438 Maven:<\/p>\n<details class=\"spoiler\">\n<summary>\u0422\u0435\u0441\u0442\u044b<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing  import java.sql.Timestamp import java.time.LocalDateTime  import org.apache.spark.sql.{SaveMode, functions => fn} import org.scalatest.flatspec.AnyFlatSpec import org.scalatest.matchers.should.Matchers._ import ru.kalininskii.orderbucketing.OrderBucketing._  import scala.util.Random  class OrderPartitionTest extends AnyFlatSpec with SparkBuilder {    val rangeLimit = 99999   val seed = 111   val numLines = 4000    private def getRandomDF(start: Int, limit: Int, days: Int, daysMinus: Int) = {     val rnd = new Random()     rnd.setSeed(seed)      val sparkSession = spark     import sparkSession.implicits._      val list = Range(start, limit)       .toList       .map(id => (Option(id), \/\/id, Option \u0434\u043b\u044f nullable = true         rnd.alphanumeric.take(rnd.nextInt(32)).mkString, \/\/name         BigDecimal(rnd.nextInt().abs * rnd.nextDouble()), \/\/amount         BigDecimal(rnd.nextLong().abs.longValue()), \/\/value         rnd.alphanumeric.take(rnd.nextInt(3)).mkString, \/\/divider         Timestamp.valueOf(LocalDateTime.now().minusSeconds((rnd.nextInt(days).abs max daysMinus) * 24 * 60 * 60)), \/\/ts         Timestamp.valueOf(LocalDateTime.now().minusSeconds((rnd.nextInt(days).abs max daysMinus) * 24 * 60 * 60)) \/\/event_time       ))      Random.shuffle(list).toDF(\"id\", \"name\", \"amount\", \"value\", \"divider\", \"ts\", \"event_time\")    }    override def beforeAll() {     super.beforeAll()     spark.sql(\"create database if not exists testschema location '\/test'\")   }    \"Range buckets partitioner\" must \"partition correctly\" in {     val dataFrameSource = getRandomDF(0, rangeLimit, 3, 0).persist()       .withColumn(\"ts_part\", fn.expr(\"substring(ts,0,10)\"))      dataFrameSource.printSchema()      val dfRep = dataFrameSource       .repartitionWithOrderAndSort(numLines, rangeLimit \/ numLines,         fn.col(\"event_time\"), List(fn.col(\"ts_part\")), List(fn.col(\"id\")))       .persist()      dfRep.explain(true)      val logicalPlanString = \"Repartition plan with unknown distribution: \" +       \"partition by ['ts_part], \" +       \"global order by 'event_time, \" +       \"local sort by ['id], \" +       \"4000 rows per task, 24 initial tasks\"      val analyzedPlanString = \"Repartition plan with unknown distribution: \" +       \"partition by [ts_part#64], \" +       \"global order by event_time#21, \" +       \"local sort by [id#15], \" +       \"4000 rows per task, 24 initial tasks\"      val sparkPlanString = \"ExchangeWithOrder Repartition with unknown distribution: \" +       \"partition by [ts_part#64], \" +       \"global order by event_time#21, \" +       \"local sort by [id#15], \" +       \"4000 rows per task, 24 initial tasks\"      assert(dfRep.queryExecution.logical.toString contains logicalPlanString)     assert(dfRep.queryExecution.analyzed.toString contains analyzedPlanString)     assert(dfRep.queryExecution.sparkPlan.toString contains sparkPlanString)      println(dfRep.rdd.getNumPartitions)      dfRep.rdd.getNumPartitions should be(28)      dfRep       .withColumn(\"rdd_part_id\", fn.spark_partition_id())       .groupBy(fn.col(\"ts_part\") as \"ts_part\", fn.col(\"rdd_part_id\"))       .agg(fn.min(\"event_time\") as \"range_min_event_time\",         fn.max(\"event_time\") as \"range_max_event_time\",         fn.count(fn.lit(0)) as \"cnt_\")       .orderBy(\"ts_part\", \"rdd_part_id\")       .show(100, false)      dfRep       .sortWithinPartitions(fn.col(\"id\"))       .write       .mode(SaveMode.Overwrite)       .partitionBy(\"ts_part\")       .format(\"parquet\")       .option(\"path\", \"\/test\/pa\/sometable\/snp\")       .saveAsTable(\"testschema.sometable\")      val dfT = spark.table(\"testschema.sometable\")      val savedCount = dfT.count     savedCount shouldEqual rangeLimit      dfT       .withColumn(\"file_name\",         fn.regexp_replace(fn.input_file_name(), \".*\/test\/pa\/sometable\/snp\", \"\"))       .groupBy(fn.col(\"ts_part\") as \"ts_part\", fn.col(\"file_name\"))       .agg(fn.min(\"event_time\") as \"range_min_event_time\",         fn.max(\"event_time\") as \"range_max_event_time\",         fn.count(fn.lit(0)) as \"cnt_\")       .orderBy(\"file_name\")       .show(100, false)      dfT.printSchema()      dfT.inputFiles.length should be(dfRep.rdd.getNumPartitions - 1)      dfT.schema should be(dataFrameSource.schema)      dfT       .withColumn(\"file_name\",         fn.regexp_replace(fn.input_file_name(), \".*\/test\/pa\/sometable\/snp\", \"\"))       .select(\"id\", \"file_name\")       .where(\"id &lt; 32\")       .show(32, false)    } }<\/code><\/pre>\n<\/div>\n<\/details>\n<details class=\"spoiler\">\n<summary>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 Hadoop Minicluster \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u043e\u0432<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing   import java.io.File import java.nio.file.Files   import org.apache.hadoop.conf.Configuration import org.apache.hadoop.fs.FileUtil import org.apache.hadoop.hdfs.MiniDFSCluster import org.apache.log4j.{Level, Logger} import org.scalatest.{BeforeAndAfterAll, TestSuite}   trait MiniHdfsBuilder extends BeforeAndAfterAll {   this: TestSuite =>     val baseDir: File = Files.createTempDirectory(\"test_hdfs\").toFile.getAbsoluteFile     val fsConf: Configuration = new Configuration()   fsConf.set(MiniDFSCluster.HDFS_MINIDFS_BASEDIR, baseDir.getAbsolutePath)     val builder: MiniDFSCluster.Builder = new MiniDFSCluster.Builder(fsConf)   var hdfsCluster: MiniDFSCluster = _     override def beforeAll() {     super.beforeAll()       Logger.getLogger(\"org.apache.hadoop\").setLevel(Level.WARN)     Logger.getLogger(\"org.apache.spark\").setLevel(Level.WARN)     Logger.getLogger(\"org.spark_project.jetty.server\").setLevel(Level.WARN)       hdfsCluster = builder.build()   }     override def afterAll() {     try {       super.afterAll()     }     finally {       hdfsCluster.shutdown()       FileUtil.fullyDelete(baseDir)     }   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<details class=\"spoiler\">\n<summary>\u0421\u043e\u0437\u0434\u0430\u043d\u0438\u0435 SparkSession \u0434\u043b\u044f \u0442\u0435\u0441\u0442\u043e\u0432<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing import java.io.File   import org.apache.hadoop.fs.FileUtil import org.apache.hadoop.hdfs.DistributedFileSystem import org.apache.spark.SparkConf import org.apache.spark.sql.SparkSession import org.scalatest.TestSuite import ru.kalininskii.orderbucketing.strategy.RepartitionStrategy   trait SparkBuilder extends MiniHdfsBuilder {   this: TestSuite =>   var fileSystem: DistributedFileSystem = _     var spark: SparkSession = _     override def beforeAll() {     super.beforeAll()     \/\/\u0441\u0435\u0441\u0441\u0438\u044f \u0443\u0436\u0435 \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u0441\u043e\u0437\u0434\u0430\u043d\u0430, \u043f\u043e\u044d\u0442\u043e\u043c\u0443 \u043d\u0443\u0436\u043d\u043e \u043f\u043e\u043f\u0440\u043e\u0431\u043e\u0432\u0430\u0442\u044c \u0435\u0435 \u043d\u0430\u0439\u0442\u0438 \u0438 \u043e\u0441\u0442\u0430\u043d\u043e\u0432\u0438\u0442\u044c     spark = SparkSession       .builder()       .appName(\"DUMMY\")       .master(\"local[*]\")       .getOrCreate()       spark.stop()       fileSystem = hdfsCluster.getFileSystem()       val sparkConf = new SparkConf().setMaster(\"local[*]\").setAppName(\"RANGE BUCKET TEST\")     sparkConf.set(\"fs.defaultFS\", fileSystem.getUri.toString)       spark = SparkSession       .builder()       .appName(\"PARTITION_TEST\")       .master(\"local[*]\")       .config(sparkConf)       .withExtensions(e => {         e.injectPlannerStrategy(_ => RepartitionStrategy)       })       .getOrCreate()       spark.sparkContext.setLogLevel(\"WARN\")   }     override def afterAll(): Unit = {     spark.stop()     FileUtil.fullyDelete(new File(\".\/metastore_db\"))     FileUtil.fullyDelete(new File(\".\/spark-warehouse\"))     FileUtil.fullyDelete(new File(\".\/derby.log\"))     FileUtil.fullyDelete(new File(\".\/chk\"))       super.afterAll()   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<details class=\"spoiler\">\n<summary>pom.xml \u0434\u043b\u044f Maven<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>&lt;?xml version=\"1.0\" encoding=\"UTF-8\"?> &lt;project xmlns=\"http:\/\/maven.apache.org\/POM\/4.0.0\"          xmlns:xsi=\"http:\/\/www.w3.org\/2001\/XMLSchema-instance\"          xsi:schemaLocation=\"http:\/\/maven.apache.org\/POM\/4.0.0 http:\/\/maven.apache.org\/xsd\/maven-4.0.0.xsd\">     &lt;modelVersion>4.0.0&lt;\/modelVersion>       &lt;groupId>fine-grained&lt;\/groupId>     &lt;artifactId>range-bucketing&lt;\/artifactId>     &lt;version>0.1-SNAPSHOT&lt;\/version>       &lt;properties>         &lt;spark.version>2.4.0&lt;\/spark.version>         &lt;hadoop.version>3.1.1&lt;\/hadoop.version>         &lt;scala.version>2.11.8&lt;\/scala.version>         &lt;maven.compiler.source>1.8&lt;\/maven.compiler.source>         &lt;maven.compiler.target>1.8&lt;\/maven.compiler.target>         &lt;project.build.sourceEncoding>UTF-8&lt;\/project.build.sourceEncoding>     &lt;\/properties>       &lt;dependencies>         &lt;dependency>             &lt;groupId>org.apache.spark&lt;\/groupId>             &lt;artifactId>spark-core_2.11&lt;\/artifactId>             &lt;version>${spark.version}&lt;\/version>             &lt;scope>provided&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.apache.spark&lt;\/groupId>             &lt;artifactId>spark-sql_2.11&lt;\/artifactId>             &lt;version>${spark.version}&lt;\/version>             &lt;scope>provided&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.scala-lang&lt;\/groupId>             &lt;artifactId>scala-library&lt;\/artifactId>             &lt;version>${scala.version}&lt;\/version>             &lt;scope>provided&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.apache.hadoop&lt;\/groupId>             &lt;artifactId>hadoop-minicluster&lt;\/artifactId>             &lt;version>${hadoop.version}&lt;\/version>             &lt;scope>test&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.apache.spark&lt;\/groupId>             &lt;artifactId>spark-hive_2.11&lt;\/artifactId>             &lt;version>${spark.version}&lt;\/version>             &lt;scope>test&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.scalatest&lt;\/groupId>             &lt;artifactId>scalatest_2.11&lt;\/artifactId>             &lt;version>3.1.1&lt;\/version>             &lt;scope>test&lt;\/scope>         &lt;\/dependency>         &lt;dependency>             &lt;groupId>org.scalactic&lt;\/groupId>             &lt;artifactId>scalactic_2.11&lt;\/artifactId>             &lt;version>3.1.1&lt;\/version>             &lt;scope>test&lt;\/scope>         &lt;\/dependency>     &lt;\/dependencies>       &lt;build>         &lt;!--   Java sources path:     -->         &lt;!--        &lt;resources>-->         &lt;!--            &lt;resource>-->         &lt;!--                &lt;directory>src\/main\/scala&lt;\/directory>-->         &lt;!--            &lt;\/resource>-->         &lt;!--        &lt;\/resources>-->         &lt;!--   For JaCoCo:     -->         &lt;sourceDirectory>src\/main\/scala&lt;\/sourceDirectory>         &lt;finalName>${project.artifactId}&lt;\/finalName>         &lt;plugins>             &lt;!-- display active profile in compile phase -->             &lt;plugin>                 &lt;groupId>org.apache.maven.plugins&lt;\/groupId>                 &lt;artifactId>maven-help-plugin&lt;\/artifactId>                 &lt;version>3.2.0&lt;\/version>                 &lt;executions>                     &lt;execution>                         &lt;id>show-profiles&lt;\/id>                         &lt;phase>compile&lt;\/phase>                         &lt;goals>                             &lt;goal>active-profiles&lt;\/goal>                         &lt;\/goals>                     &lt;\/execution>                 &lt;\/executions>             &lt;\/plugin>             &lt;plugin>                 &lt;groupId>org.apache.maven.plugins&lt;\/groupId>                 &lt;artifactId>maven-jar-plugin&lt;\/artifactId>                 &lt;version>3.2.0&lt;\/version>                 &lt;configuration>                     &lt;archive>                         &lt;manifest>                             &lt;addDefaultImplementationEntries>true&lt;\/addDefaultImplementationEntries>                             &lt;addDefaultSpecificationEntries>true&lt;\/addDefaultSpecificationEntries>                         &lt;\/manifest>                     &lt;\/archive>                 &lt;\/configuration>             &lt;\/plugin>             &lt;!-- scala and java mix compilation -->             &lt;plugin>                 &lt;groupId>net.alchim31.maven&lt;\/groupId>                 &lt;artifactId>scala-maven-plugin&lt;\/artifactId>                 &lt;version>3.2.2&lt;\/version>                 &lt;executions>                     &lt;execution>                         &lt;goals>                             &lt;goal>compile&lt;\/goal>                             &lt;goal>testCompile&lt;\/goal>                         &lt;\/goals>                     &lt;\/execution>                 &lt;\/executions>                 &lt;configuration>                     &lt;scalaVersion>${scala.version}&lt;\/scalaVersion>                 &lt;\/configuration>             &lt;\/plugin>               &lt;plugin>                 &lt;groupId>org.apache.maven.plugins&lt;\/groupId>                 &lt;artifactId>maven-surefire-plugin&lt;\/artifactId>                 &lt;version>2.7&lt;\/version>                 &lt;configuration>                     &lt;skipTests>true&lt;\/skipTests>                 &lt;\/configuration>             &lt;\/plugin>             &lt;plugin>                 &lt;groupId>org.scalatest&lt;\/groupId>                 &lt;artifactId>scalatest-maven-plugin&lt;\/artifactId>                 &lt;version>1.0&lt;\/version>                 &lt;executions>                     &lt;execution>                         &lt;id>test&lt;\/id>                         &lt;goals>                             &lt;goal>test&lt;\/goal>                         &lt;\/goals>                     &lt;\/execution>                 &lt;\/executions>             &lt;\/plugin>             &lt;plugin>                 &lt;groupId>org.jacoco&lt;\/groupId>                 &lt;artifactId>jacoco-maven-plugin&lt;\/artifactId>                 &lt;version>0.8.5&lt;\/version>                 &lt;executions>                     &lt;execution>                         &lt;goals>                             &lt;goal>prepare-agent&lt;\/goal>                         &lt;\/goals>                     &lt;\/execution>                     &lt;!-- attached to Maven test phase -->                     &lt;execution>                         &lt;id>report&lt;\/id>                         &lt;phase>test&lt;\/phase>                         &lt;goals>                             &lt;goal>report&lt;\/goal>                         &lt;\/goals>                     &lt;\/execution>                 &lt;\/executions>             &lt;\/plugin>         &lt;\/plugins>     &lt;\/build>   &lt;\/project><\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041d\u0430 \u044d\u0442\u043e\u043c \u0441\u0442\u0430\u0442\u044c\u044f \u0437\u0430\u043a\u043e\u043d\u0447\u0435\u043d\u0430, \u0430 \u043d\u0430\u0448\u0430 \u0440\u0430\u0431\u043e\u0442\u0430 \u043f\u0440\u043e\u0434\u043e\u043b\u0436\u0430\u0435\u0442\u0441\u044f. \u0420\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u043a\u0430 \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d\u0438\u044f Spark \u0442\u0440\u0435\u0431\u0443\u0435\u0442 \u043c\u043d\u043e\u0436\u0435\u0441\u0442\u0432\u0430 \u0440\u0435\u0448\u0435\u043d\u0438\u0439, \u0438 \u043d\u0435 \u0432\u0441\u0435\u0433\u0434\u0430 \u043e\u043d\u0438 \u043f\u0440\u043e\u0441\u0442\u044b\u0435. \u042d\u0442\u043e \u0432\u0438\u0434\u043d\u043e \u0438 \u0432 \u0438\u0441\u0445\u043e\u0434\u043d\u043e\u043c \u043a\u043e\u0434\u0435 Spark, \u044f \u0434\u0443\u043c\u0430\u044e, \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0447\u0438\u043a\u0438 \u0441\u0435\u0439\u0447\u0430\u0441 \u0440\u0435\u0430\u043b\u0438\u0437\u043e\u0432\u0430\u043b\u0438 \u0431\u044b \u043c\u043d\u043e\u0433\u043e\u0435 \u0438\u043d\u0430\u0447\u0435.<\/p>\n<p>\u0412 \u043b\u044e\u0431\u043e\u043c \u0441\u043b\u0443\u0447\u0430\u0435, \u043c\u044b \u043d\u0435 \u043e\u0431\u044f\u0437\u0430\u043d\u044b \u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c\u0441\u044f \u0442\u043e\u043b\u044c\u043a\u043e \u0441\u0443\u0449\u0435\u0441\u0442\u0432\u0443\u044e\u0449\u0438\u043c\u0438 \u0441\u0440\u0435\u0434\u0441\u0442\u0432\u0430\u043c\u0438, \u043c\u043e\u0436\u043d\u043e \u0440\u0430\u0437\u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u0447\u0442\u043e-\u0442\u043e \u0441\u0432\u043e\u0451. \u0415\u0441\u043b\u0438 \u0432\u044b \u0445\u043e\u0442\u0438\u0442\u0435 \u043f\u0440\u0435\u043e\u0431\u0440\u0430\u0437\u043e\u0432\u0430\u0442\u044c \u0434\u0430\u043d\u043d\u044b\u0435, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f Spark, \u0442\u043e \u043c\u043e\u0436\u0435\u0442\u0435 \u0432\u044b\u0431\u0440\u0430\u0442\u044c \u0433\u043e\u0442\u043e\u0432\u043e\u0435 \u0440\u0435\u0448\u0435\u043d\u0438\u0435 (HashPartitioner, RangePartitioner). \u0418 \u0432\u0441\u0451 \u0436\u0435, \u0435\u0441\u043b\u0438 \u0435\u0441\u043b\u0438 \u0432\u044b \u043f\u043e\u043d\u044f\u043b\u0438, \u0447\u0442\u043e \u043d\u0443\u0436\u043d\u043e \u0441\u043f\u0435\u0446\u0438\u0444\u0438\u0447\u0435\u0441\u043a\u043e\u0435 \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0435 \u0434\u0430\u043d\u043d\u044b\u0445, \u0442\u043e \u043d\u0435 \u0431\u043e\u0439\u0442\u0435\u0441\u044c \u0441\u0434\u0435\u043b\u0430\u0442\u044c \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u0434\u043b\u044f \u0441\u0432\u043e\u0438\u0445 \u043d\u0443\u0436\u0434, \u0438 \u044f \u0443\u0432\u0435\u0440\u0435\u043d, \u0447\u0442\u043e \u043e\u043d \u0431\u0443\u0434\u0435\u0442 \u043b\u0443\u0447\u0448\u0435 \u0438 \u043f\u043e\u043b\u0435\u0437\u043d\u0435\u0435!<\/p>\n<p>\u041d\u0435 \u043f\u044b\u0442\u0430\u0439\u0442\u0435\u0441\u044c \u0437\u0430\u0441\u0442\u0430\u0432\u0438\u0442\u044c \u043c\u0430\u0448\u0438\u043d\u0443 \u0440\u0430\u0431\u043e\u0442\u0430\u0442\u044c \u0431\u044b\u0441\u0442\u0440\u0435\u0435, \u043f\u043e\u0441\u0442\u0430\u0440\u0430\u0439\u0442\u0435\u0441\u044c \u0441\u0434\u0435\u043b\u0430\u0442\u044c \u0442\u0430\u043a, \u0447\u0442\u043e\u0431\u044b \u043c\u0430\u0448\u0438\u043d\u0430 \u043d\u0435 \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u043b\u0430 \u043d\u0435\u043d\u0443\u0436\u043d\u0443\u044e \u0440\u0430\u0431\u043e\u0442\u0443, \u0430 \u043d\u0443\u0436\u043d\u043e\u0439 \u0440\u0430\u0431\u043e\u0442\u044b \u0431\u044b\u043b\u043e \u043a\u0430\u043a \u043c\u043e\u0436\u043d\u043e \u043c\u0435\u043d\u044c\u0448\u0435. \u0422\u043e\u0433\u0434\u0430 \u044d\u0442\u043e \u0431\u0443\u0434\u0435\u0442 \u043d\u0435 \u043f\u0440\u0435\u0436\u0434\u0435\u0432\u0440\u0435\u043c\u0435\u043d\u043d\u0430\u044f, \u0430 \u043e\u0441\u043e\u0437\u043d\u0430\u043d\u043d\u0430\u044f \u0438 \u043d\u0443\u0436\u043d\u0430\u044f \u043e\u043f\u0442\u0438\u043c\u0438\u0437\u0430\u0446\u0438\u044f.<\/p>\n<p>\u0412 \u0441\u043b\u0435\u0434\u0443\u044e\u0449\u0435\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u043c\u044b \u0440\u0430\u0437\u0431\u0435\u0440\u0451\u043c \u0441\u0430\u043c DataSource, \u043f\u0440\u0435\u0434\u043d\u0430\u0437\u043d\u0430\u0447\u0435\u043d\u043d\u044b\u0439 \u0434\u043b\u044f \u0447\u0442\u0435\u043d\u0438\u044f \u0434\u0430\u043d\u043d\u044b\u0445, \u0438 \u0441\u0440\u0435\u0434\u0441\u0442\u0432\u0430 \u0440\u0430\u0431\u043e\u0442\u044b \u0441 \u043c\u0435\u0442\u0430\u0434\u0430\u043d\u043d\u044b\u043c\u0438.<\/p>\n<\/div>\n<\/div>\n<\/div>\n<p><!----><!----><\/div>\n<p><!----><!----><br \/> \u0441\u0441\u044b\u043b\u043a\u0430 \u043d\u0430 \u043e\u0440\u0438\u0433\u0438\u043d\u0430\u043b \u0441\u0442\u0430\u0442\u044c\u0438 <a href=\"https:\/\/habr.com\/ru\/articles\/583018\/\"> https:\/\/habr.com\/ru\/articles\/583018\/<\/a><\/p>\n","protected":false},"excerpt":{"rendered":"<div><!--[--><!--]--><\/div>\n<div id=\"post-content-body\">\n<div>\n<div class=\"article-formatted-body article-formatted-body article-formatted-body_version-2\">\n<div xmlns=\"http:\/\/www.w3.org\/1999\/xhtml\">\n<p><strong><em>\u0410\u0432\u0442\u043e\u0440:<\/em><\/strong><em> \u0418\u0432\u0430\u043d \u041a\u0430\u043b\u0438\u043d\u0438\u043d\u0441\u043a\u0438\u0439, \u0443\u0447\u0430\u0441\u0442\u043d\u0438\u043a \u043f\u0440\u043e\u0444\u0435\u0441\u0441\u0438\u043e\u043d\u0430\u043b\u044c\u043d\u043e\u0433\u043e \u0441\u043e\u043e\u0431\u0449\u0435\u0441\u0442\u0432\u0430 \u0421\u0431\u0435\u0440\u0430 SberProfi DWH\/BigData.<\/em><\/p>\n<p><em>\u041f\u0440\u043e\u0444\u0435\u0441\u0441\u0438\u043e\u043d\u0430\u043b\u044c\u043d\u043e\u0435 \u0441\u043e\u043e\u0431\u0449\u0435\u0441\u0442\u0432\u043e SberProfi DWH\/BigData \u043e\u0442\u0432\u0435\u0447\u0430\u0435\u0442 \u0437\u0430 \u0440\u0430\u0437\u0432\u0438\u0442\u0438\u0435 \u043a\u043e\u043c\u043f\u0435\u0442\u0435\u043d\u0446\u0438\u0439 \u0432 \u0442\u0430\u043a\u0438\u0445 \u043d\u0430\u043f\u0440\u0430\u0432\u043b\u0435\u043d\u0438\u044f\u0445, \u043a\u0430\u043a \u044d\u043a\u043e\u0441\u0438\u0441\u0442\u0435\u043c\u0430 Hadoop, Teradata, Oracle DB, GreenPlum, \u0430 \u0442\u0430\u043a\u0436\u0435 BI \u0438\u043d\u0441\u0442\u0440\u0443\u043c\u0435\u043d\u0442\u0430\u0445 Qlik, SAP BO, Tableau \u0438 \u0434\u0440.<\/em><\/p>\n<p>\u041d\u0430\u0447\u043d\u0443 \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u0435 \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440\u0430 \u0441 UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u044b:<\/p>\n<figure class=\"full-width\"><figcaption>\u0420\u0438\u0441\u0443\u043d\u043e\u043a 1. UML-\u0434\u0438\u0430\u0433\u0440\u0430\u043c\u043c\u0430 \u043a\u043b\u0430\u0441\u0441\u043e\u0432 OrderBucketsPartitioner<\/figcaption><\/figure>\n<p>\u041d\u0430\u0432\u0435\u0440\u043d\u044f\u043a\u0430 \u0432\u044b \u0432\u0438\u0434\u0435\u043b\u0438, \u043a\u0430\u043a Spark \u0440\u0430\u0437\u0431\u0438\u0440\u0430\u0435\u0442 \u0437\u0430\u043f\u0440\u043e\u0441, \u0441\u0442\u0440\u043e\u0438\u0442 \u043f\u043b\u0430\u043d \u0438 \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0435\u0442 \u0435\u0433\u043e. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u0438\u043c \u044d\u0442\u0443 \u0441\u0445\u0435\u043c\u0443 \u0441 \u043d\u0430\u0448\u0435\u0439 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0435\u0439 (\u043a\u043e\u0433\u0434\u0430 \u0431\u0443\u0434\u0435\u0442\u0435 \u0437\u043d\u0430\u043a\u043e\u043c\u0438\u0442\u044c\u0441\u044f \u0441 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0435\u0439 \u043a\u043b\u0430\u0441\u0441\u043e\u0432, \u0432\u043e\u0437\u0432\u0440\u0430\u0449\u0430\u0439\u0442\u0435\u0441\u044c \u043a \u044d\u0442\u043e\u0439 \u0441\u0445\u0435\u043c\u0435, \u0447\u0442\u043e\u0431\u044b \u0443\u0432\u0438\u0434\u0435\u0442\u044c \u0441\u043e\u043e\u0442\u0432\u0435\u0442\u0441\u0442\u0432\u0438\u0435):<\/p>\n<figure class=\"full-width\"><figcaption>\u0420\u0438\u0441\u0443\u043d\u043e\u043a 2. \u0421\u043e\u043f\u043e\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u0435 \u0441\u0445\u0435\u043c\u044b \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0438 \u0437\u0430\u043f\u0440\u043e\u0441\u0430 \u0438 \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0439 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438.<\/figcaption><\/figure>\n<p>\u041d\u0443\u0436\u0435\u043d \u043c\u0435\u0442\u043e\u0434, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u0437\u0432\u043e\u043b\u0438\u043b \u0431\u044b \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u0434\u043b\u044f \u043b\u044e\u0431\u043e\u0433\u043e \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430. \u041f\u043e\u043a\u0430 \u043d\u0435 \u0431\u0443\u0434\u0435\u043c \u0441\u0432\u044f\u0437\u044b\u0432\u0430\u0442\u044c\u0441\u044f \u0441 SQL (\u0440\u0430\u0437\u0432\u0435 \u043a\u0442\u043e-\u0442\u043e \u0434\u0435\u043b\u0430\u0435\u0442 dataframe.repartition(&#8230;) \u0441 \u043f\u043e\u043c\u043e\u0449\u044c\u044e SQL?).<\/p>\n<p>\u0414\u043e\u0431\u0430\u0432\u0438\u043c \u0442\u0430\u043a\u043e\u0439 \u043c\u0435\u0442\u043e\u0434 \u0432 Scala Dataframe API, \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u044f \u043f\u0430\u0442\u0442\u0435\u0440\u043d \u00abPimp my library\u00bb.<br \/> \u0414\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0441\u043e\u0437\u0434\u0430\u0434\u0438\u043c \u043d\u043e\u0432\u044b\u0439 \u043e\u0431\u044a\u0435\u043a\u0442 ru.kalininskii.orderbucketing.OrderBucketing. \u041e\u043d \u0431\u0443\u0434\u0435\u0442 \u0441\u043e\u0434\u0435\u0440\u0436\u0430\u0442\u044c implicit class \u0441 \u043f\u043e\u043a\u0430 \u0435\u0434\u0438\u043d\u0441\u0442\u0432\u0435\u043d\u043d\u044b\u043c \u043c\u0435\u0442\u043e\u0434\u043e\u043c repartitionWithOrderAndSort:<\/p>\n<details class=\"spoiler\">\n<summary>OrderBucketing<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing   import org.apache.spark.sql.{Column, DataFrame, PlanHelper} import org.apache.spark.sql.catalyst.expressions.{Ascending, Expression, SortOrder} import ru.kalininskii.orderbucketing.plans.logical.RepartitionWithOrderAndSort   object OrderBucketing {     implicit class DataFramePartOps(ds: DataFrame) {     def repartitionWithOrderAndSort(numLines: Int,                                      numPartitions: Int,                                      orderColumn: Column,                                      partitionColumns: Seq[Column],                                      sortColumns: Seq[Column]): DataFrame = {       def toSortOrder(col: Column): SortOrder = {         col.expr match {           case order: SortOrder => order.copy(direction = Ascending)           case expr: Expression => SortOrder(expr, Ascending)           case _ => throw new Exception(s\"Can not get order from $col\")         }       }         val orderExpression = toSortOrder(orderColumn)       val partitionExpressions = partitionColumns.map(_.expr)       val sortExpressions = sortColumns.map(toSortOrder)         val logicalPlan = RepartitionWithOrderAndSort(         orderExpression,         partitionExpressions,         sortExpressions,         numLines,         numPartitions,         None,         ds.queryExecution.logical)         PlanHelper.planToDF(ds.sparkSession, logicalPlan)     }   }   }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041a\u0430\u043a \u043c\u043e\u0436\u043d\u043e \u0432\u0438\u0434\u0435\u0442\u044c \u0438\u0437 \u043a\u043e\u0434\u0430 \u0432\u044b\u0448\u0435, \u043c\u0435\u0442\u043e\u0434 \u043f\u0440\u0438\u043d\u0438\u043c\u0430\u0435\u0442 \u0430\u0440\u0433\u0443\u043c\u0435\u043d\u0442\u00a0partitionColumns: Seq[Column], \u0432 \u043a\u043e\u0442\u043e\u0440\u043e\u043c \u043f\u043e\u043b\u044f \u044f\u0432\u043d\u043e\u0433\u043e \u0438 \u043d\u0435\u044f\u0432\u043d\u043e\u0433\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043e\u0431\u044a\u0435\u0434\u0438\u043d\u0435\u043d\u044b. \u0422\u0430\u043a \u0438 \u0434\u043e\u043b\u0436\u043d\u043e \u0431\u044b\u0442\u044c, \u043f\u0430\u0440\u0442\u0438\u0448\u0435\u043d\u0435\u0440 \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u0435\u0442 \u0437\u0430\u043f\u0438\u0441\u0438 \u043f\u043e \u0441\u0435\u043a\u0446\u0438\u044f\u043c \u043d\u0435\u0437\u0430\u0432\u0438\u0441\u0438\u043c\u043e \u043e\u0442 \u0438\u0445 \u0432\u043d\u0435\u0448\u043d\u0435\u0433\u043e \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u043b\u0435\u043d\u0438\u044f.\u00a0<\/p>\n<p>\u041e\u0431\u044a\u0435\u043a\u0442 PlanHelper \u043f\u043e\u043a\u0430 \u0447\u0442\u043e \u0442\u0430\u043a\u0436\u0435 \u043e\u0431\u043e\u0439\u0434\u0451\u0442\u0441\u044f \u043e\u0434\u043d\u0438\u043c \u043c\u0435\u0442\u043e\u0434\u043e\u043c, planToDF. \u042d\u0442\u043e\u0442 \u043c\u0435\u0442\u043e\u0434 \u0432\u044b\u0437\u044b\u0432\u0430\u0435\u0442 \u043f\u0440\u0438\u0432\u0430\u0442\u043d\u044b\u0439 \u043c\u0435\u0442\u043e\u0434 \u0434\u043b\u044f \u043f\u0430\u043a\u0435\u0442\u0430 org.apache.spark.sql, \u043f\u043e\u044d\u0442\u043e\u043c\u0443, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u043c\u043e\u0433 \u0432\u044b\u043f\u043e\u043b\u043d\u044f\u0442\u044c\u0441\u044f, \u043e\u0431\u044a\u0435\u043a\u0442 \u0442\u043e\u0436\u0435 \u0434\u043e\u043b\u0436\u0435\u043d \u043d\u0430\u0445\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0432 \u044d\u0442\u043e\u043c \u043f\u0430\u043a\u0435\u0442\u0435:<\/p>\n<details class=\"spoiler\">\n<summary>PlanHelper<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package org.apache.spark.sql   import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan   object PlanHelper {     def planToDF(spark: SparkSession, logicalPlan: LogicalPlan): DataFrame = {     Dataset.ofRows(spark, logicalPlan)   } }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0438\u043f \u0434\u043b\u044f \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD, \u0430 \u0442\u0430\u043a\u0436\u0435 \u043a\u043b\u044e\u0447\u0430 RDD, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043f\u043e\u043d\u0430\u0434\u043e\u0431\u0438\u0442\u0441\u044f \u043d\u0430\u043c \u0432 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0451\u043d\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442, \u043f\u043e\u043c\u0435\u0441\u0442\u0438\u043c \u0432 package object, \u0447\u0442\u043e\u0431\u044b \u043e\u043d \u0431\u044b\u043b \u0434\u043e\u0441\u0442\u0443\u043f\u0435\u043d \u0432\u0441\u0435\u043c\u0443 \u043d\u0430\u0448\u0435\u043c\u0443 \u043f\u0430\u043a\u0435\u0442\u0443 \u0438 \u043d\u0435 \u0442\u043e\u043b\u044c\u043a\u043e<\/p>\n<details class=\"spoiler\">\n<summary>package object rangebucketing<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii   import org.apache.spark.sql.catalyst.InternalRow   package object rangebucketing {   type BucketsDistribution = (Int, InternalRow, Either[(InternalRow, Seq[String]), (InternalRow, InternalRow, Int)])   type OrderAndSortKey = (InternalRow, InternalRow, InternalRow) }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u0422\u0435\u043f\u0435\u0440\u044c \u043d\u0443\u0436\u043d\u043e \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u043f\u043b\u0430\u043d\u044b, \u043d\u0430\u0447\u0438\u043d\u0430\u044f \u0441 \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e. \u0418 \u0432 \u043d\u0438\u0445 \u043f\u0435\u0440\u0435\u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0438\u0442\u044c \u043c\u0435\u0442\u043e\u0434 <em>simpleString<\/em>, \u0442\u0430\u043a \u043a\u0430\u043a \u0438\u043d\u0430\u0447\u0435 \u043e\u043d \u0432\u044b\u0432\u0435\u0434\u0435\u0442 \u0432 \u043f\u043b\u0430\u043d \u0432\u0435\u0441\u044c \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0439 \u043e\u0431\u044a\u0435\u043a\u0442 <em>distribution.\u00a0<\/em>\u041b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u0438\u0439 \u043f\u043b\u0430\u043d \u0434\u043b\u044f \u044d\u0442\u043e\u0433\u043e \u0440\u0435\u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u0431\u0443\u0434\u0435\u0442 \u0432\u0438\u0434\u0435\u043d \u043f\u0440\u0438 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u043d\u0438\u0438 \u043c\u0435\u0442\u043e\u0434\u0430 explain, \u0438 \u043c\u043e\u0436\u043d\u043e \u0431\u0443\u0434\u0435\u0442 \u0443\u0431\u0435\u0434\u0438\u0442\u044c\u0441\u044f, \u0447\u0442\u043e \u043e\u043d \u0434\u0435\u0439\u0441\u0442\u0432\u0438\u0442\u0435\u043b\u044c\u043d\u043e \u043f\u0440\u0438\u043c\u0435\u043d\u044f\u0435\u0442\u0441\u044f. \u041a\u0440\u043e\u043c\u0435 \u0442\u043e\u0433\u043e, \u043e\u043d \u0432\u0441\u0442\u0440\u043e\u0435\u043d \u0432 \u0444\u043e\u0440\u043c\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0435 \u043f\u043b\u0430\u043d\u043e\u0432 \u0432\u044b\u043f\u043e\u043b\u043d\u0435\u043d\u0438\u044f \u0437\u0430\u043f\u0440\u043e\u0441\u043e\u0432 \u0438 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u0435\u0434\u043e\u0441\u0442\u0430\u0432\u043b\u044f\u0442\u044c \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044e \u043e \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u043e\u043c \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u0438 \u043d\u0430\u0431\u043e\u0440\u0430 \u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u043f\u0435\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 <em>partitioning<\/em><strong><em>.<\/em><\/strong><\/p>\n<details class=\"spoiler\">\n<summary>RepartitionWithOrderAndSort<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing.plans.logical   import org.apache.spark.sql.catalyst.expressions.{Expression, SortOrder} import org.apache.spark.sql.catalyst.plans.logical.{LogicalPlan, RepartitionOperation} import ru.kalininskii.orderbucketing.BucketsDistribution import ru.kalininskii.orderbucketing.plans.physical.OrderBucketsPartitioning   \/**  * \u042d\u0442\u043e\u0442 \u043a\u043b\u0430\u0441\u0441 \u0440\u0430\u0437\u0434\u0435\u043b\u044f\u0435\u0442 \u0434\u0430\u043d\u043d\u044b\u0435 \u043f\u043e \u0443\u0441\u043b\u043e\u0432\u044e \u0443\u043d\u0438\u043a\u0430\u043b\u044c\u043d\u044b\u0445 \u0442\u043e\u0447\u043d\u044b\u0445 \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 [[Expression]]  *  * @param orderExpression      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u0435, \u043f\u043e \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043a\u043e\u0442\u043e\u0440\u043e\u0433\u043e \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  *                             \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f [[partitionExpressions]]  * @param partitionExpressions \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435.  * @param sortExpressions      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u0430\u044f \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0430.  * @param numLines             \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD (\u0438 \u0432 \u0437\u0430\u043f\u0438\u0441\u0430\u043d\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435)  * @param numPartitions        \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439, \u043c\u043e\u0436\u0435\u0442 \u043e\u0442\u043b\u0438\u0447\u0430\u0442\u044c\u0441\u044f  * @param distribution         \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044f \u043e\u0431 \u0438\u043c\u0435\u044e\u0449\u0435\u043c\u0441\u044f \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0438, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043d\u0430\u0434\u043e \u0432\u043e\u0441\u043f\u0440\u043e\u0438\u0437\u0432\u0435\u0441\u0442\u0438  * @param child                \u043f\u043b\u0430\u043d \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u0438\u044f \u043d\u0430\u0431\u043e\u0440\u0430 \u0434\u0430\u043d\u043d\u044b\u0445, \u043a\u043e\u0442\u043e\u0440\u044b\u0439 \u043d\u0443\u0436\u043d\u043e \u0441\u0435\u043a\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u0442\u044c  *\/ case class RepartitionWithOrderAndSort(                                         orderExpression: SortOrder,                                         partitionExpressions: Seq[Expression],                                         sortExpressions: Seq[SortOrder],                                         numLines: Int,                                         numPartitions: Int,                                         distribution: Option[Seq[BucketsDistribution]],                                         child: LogicalPlan                                       ) extends RepartitionOperation {     override def nodeName: String = s\"Repartition plan with \" +     s\"${distribution.map(_ => \"predefined\").getOrElse(\"unknown\")} distribution:\"     override def simpleString: String = {     s\"$nodeName ${partitionExpressions.mkString(\"partition by [\", \", \", \"], \")}\" +       s\"global order by ${orderExpression.child}, \" +       s\"${sortExpressions.map(_.child).mkString(\"local sort by [\", \", \", \"], \")}\" +       numLines + \" rows per task, \" + numPartitions + \" initial tasks\"   }     val partitioning: OrderBucketsPartitioning = OrderBucketsPartitioning(     orderExpression, partitionExpressions, sortExpressions, numLines, numPartitions, distribution)     override def maxRows: Option[Long] = child.maxRows     override def shuffle: Boolean = true }<\/code><\/pre>\n<\/div>\n<\/details>\n<p>\u041a\u043b\u0430\u0441\u0441 <em>Partitioning <\/em>\u0438 \u0435\u0433\u043e \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d\u0438\u044f\u00a0\u0432 \u0441\u0442\u0440\u0443\u043a\u0442\u0443\u0440\u0435 \u043f\u0430\u043a\u0435\u0442\u043e\u0432 Spark \u043e\u0442\u043d\u043e\u0441\u044f\u0442\u0441\u044f \u043a \u0444\u0438\u0437\u0438\u0447\u0435\u0441\u043a\u0438\u043c \u043f\u043b\u0430\u043d\u0430\u043c (<em>org.apache.spark.sql.catalyst.plans.physical<\/em>), \u043d\u043e \u0441\u043a\u043e\u0440\u0435\u0435 \u043f\u0440\u0435\u0434\u0441\u0442\u0430\u0432\u043b\u044f\u044e\u0442 \u0441\u043e\u0431\u043e\u0439 \u043a\u043e\u043d\u0442\u0435\u0439\u043d\u0435\u0440\u044b \u0434\u043b\u044f \u0445\u0430\u0440\u0430\u043a\u0442\u0435\u0440\u0438\u0441\u0442\u0438\u043a \u043a\u043e\u043d\u043a\u0440\u0435\u0442\u043d\u043e\u0433\u043e \u0434\u0430\u0442\u0430\u0444\u0440\u0435\u0439\u043c\u0430.<\/p>\n<p>\u0412\u0441\u0435 \u0432\u0441\u0442\u0440\u043e\u0435\u043d\u043d\u044b\u0435 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0441\u043e\u0437\u0434\u0430\u044e\u0442\u0441\u044f \u043a\u0435\u0439\u0441 \u043a\u043b\u0430\u0441\u0441\u043e\u043c (\u0441\u043c. \u0422\u0435\u0440\u043c\u0438\u043d\u044b \u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f)\u00a0<em>org.apache.spark.sql.catalyst.plans.logical.RepartitionByExpression,\u00a0<\/em>\u043f\u043e\u043b\u0443\u0447\u0430\u044f \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f \u043f\u043e\u043b\u0435\u0439, \u043f\u0435\u0440\u0435\u0434\u0430\u043d\u043d\u044b\u0445 \u0432 \u0443\u043a\u0430\u0437\u0430\u043d\u043d\u044b\u0439 \u043a\u0435\u0439\u0441 \u043a\u043b\u0430\u0441\u0441.<\/p>\n<p>\u0418 \u0432 \u043d\u0430\u0448\u0435\u043c \u0441\u043b\u0443\u0447\u0430\u0435,\u00a0<em>OrderBucketsPartitioning<\/em>\u00a0\u043c\u043e\u0436\u043d\u043e \u0441\u0447\u0438\u0442\u0430\u0442\u044c \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u0435\u043c \u043e\u043f\u0435\u0440\u0430\u0446\u0438\u0438 \u043d\u0430\u0434 \u0442\u0430\u0431\u043b\u0438\u0446\u0435\u0439, \u043f\u0440\u043e\u0447\u0438\u0442\u0430\u043d\u043d\u043e\u0439 \u0438\u0437 \u0444\u0430\u0439\u043b\u043e\u0432\u043e\u0439 \u0441\u0438\u0441\u0442\u0435\u043c\u044b \u0438\u043b\u0438 \u043e\u0431\u044a\u0435\u043a\u0442\u043d\u043e\u0433\u043e \u0445\u0440\u0430\u043d\u0438\u043b\u0438\u0449\u0430. \u041d\u0435 \u0431\u0443\u0434\u0435\u0442 \u043e\u0448\u0438\u0431\u043a\u043e\u0439 \u043e\u0442\u043d\u0435\u0441\u0442\u0438 \u043a\u043b\u0430\u0441\u0441 \u043a \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u043c\u0443 \u043f\u043b\u0430\u043d\u0443, \u0432\u0435\u0434\u044c \u043c\u0435\u0442\u043e\u0434 <em>dataframe.explain(true)<\/em> \u0438\u043c\u0435\u043d\u043d\u043e \u0435\u0433\u043e \u0432\u044b\u0432\u043e\u0434\u0438\u0442 \u0432 \u043a\u0430\u0436\u0434\u043e\u043c \u0434\u0435\u0440\u0435\u0432\u0435 \u043b\u043e\u0433\u0438\u0447\u0435\u0441\u043a\u043e\u0433\u043e \u043f\u043b\u0430\u043d\u0430).<\/p>\n<p>\u0412 \u0440\u0435\u0430\u043b\u0438\u0437\u0430\u0446\u0438\u0438 \u0435\u0441\u0442\u044c \u043d\u0435\u043f\u0440\u0438\u044f\u0442\u043d\u044b\u0439 \u043c\u043e\u043c\u0435\u043d\u0442: \u0442\u0440\u0435\u0439\u0442 (\u0441\u043c. \u0422\u0435\u0440\u043c\u0438\u043d\u044b \u0438 \u043e\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f) <em>Distribution \u00a0 <\/em>\u0432 \u0438\u0441\u0445\u043e\u0434\u043d\u043e\u043c \u043a\u043e\u0434\u0435 Spark \u043e\u0431\u044a\u044f\u0432\u043b\u0435\u043d \u043a\u0430\u043a <em>sealed<\/em>, \u0430 \u0437\u043d\u0430\u0447\u0438\u0442 \u043d\u0435 \u043c\u043e\u0436\u0435\u0442 \u0431\u044b\u0442\u044c \u0440\u0430\u0441\u0448\u0438\u0440\u0435\u043d. \u041f\u043e\u044d\u0442\u043e\u043c\u0443 \u0432 \u043c\u0435\u0442\u043e\u0434\u0435 <em>satisfies0<\/em> \u044f \u0431\u0443\u0434\u0443 \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u043e\u0432\u0430\u0442\u044c <em>OrderedDistribution<\/em>, \u0445\u043e\u0442\u044f \u043e\u043d\u0430 \u043d\u0435 \u043b\u0443\u0447\u0448\u0438\u043c \u043e\u0431\u0440\u0430\u0437\u043e\u043c \u043f\u043e\u0434\u0445\u043e\u0434\u0438\u0442 \u0434\u043b\u044f \u043e\u043f\u0438\u0441\u0430\u043d\u0438\u044f \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u043e\u0433\u043e \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u044f.<\/p>\n<details class=\"spoiler\">\n<summary>OrderBucketsPartitioning<\/summary>\n<div class=\"spoiler__content\">\n<pre><code>package ru.kalininskii.orderbucketing.plans.physical   import org.apache.spark.sql.catalyst.expressions.{Expression, SortOrder, Unevaluable} import org.apache.spark.sql.catalyst.plans.physical.{Distribution, OrderedDistribution, Partitioning} import org.apache.spark.sql.types.{DataType, IntegerType} import ru.kalininskii.orderbucketing.BucketsDistribution   \/**  * \u042d\u0442\u043e\u0442 \u043a\u043b\u0430\u0441\u0441 \u043f\u0435\u0440\u0435\u0434\u0430\u0451\u0442 \u0434\u0430\u043d\u043d\u044b\u0435 \u0434\u043b\u044f \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u043e\u043d\u0438\u0440\u043e\u0432\u0430\u043d\u0438\u044f \u043f\u043e \u0434\u0438\u0430\u043f\u0430\u0437\u043e\u043d\u0430\u043c \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043f\u0430\u0440\u0442\u0438\u0446\u0438\u0439  *  * @param orderExpression      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u0435, \u043f\u043e \u0438\u043d\u0442\u0435\u0440\u0432\u0430\u043b\u0430\u043c \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u0439 \u043a\u043e\u0442\u043e\u0440\u043e\u0433\u043e \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  *                             \u0432 \u043f\u0440\u0435\u0434\u0435\u043b\u0430\u0445 \u043e\u0434\u043d\u043e\u0433\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f [[partitionExpressions]]  * @param partitionExpressions \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u0440\u0430\u0437\u0434\u0435\u043b\u0435\u043d\u0438\u0435  * @param sortExpressions      \u0432\u044b\u0440\u0430\u0436\u0435\u043d\u0438\u044f, \u043f\u043e \u0437\u043d\u0430\u0447\u0435\u043d\u0438\u044f\u043c \u043a\u043e\u0442\u043e\u0440\u044b\u0445 \u0431\u0443\u0434\u0435\u0442 \u043f\u0440\u043e\u0438\u0437\u0432\u043e\u0434\u0438\u0442\u044c\u0441\u044f \u043b\u043e\u043a\u0430\u043b\u044c\u043d\u0430\u044f \u0441\u043e\u0440\u0442\u0438\u0440\u043e\u0432\u043a\u0430  * @param numLines             \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0437\u0430\u043f\u0438\u0441\u0435\u0439 \u0432 \u043e\u0434\u043d\u043e\u0439 \u0441\u0435\u043a\u0446\u0438\u0438 RDD (\u0438 \u0432 \u0437\u0430\u043f\u0438\u0441\u0430\u043d\u043d\u043e\u043c \u0444\u0430\u0439\u043b\u0435)  * @param numPartitions        \u043f\u0440\u0435\u0434\u043f\u043e\u043b\u0430\u0433\u0430\u0435\u043c\u043e\u0435 \u043a\u043e\u043b\u0438\u0447\u0435\u0441\u0442\u0432\u043e \u0441\u0435\u043a\u0446\u0438\u0439, \u043c\u043e\u0436\u0435\u0442 \u043e\u0442\u043b\u0438\u0447\u0430\u0442\u044c\u0441\u044f  * @param distribution         \u0438\u043d\u0444\u043e\u0440\u043c\u0430\u0446\u0438\u044f \u043e\u0431 \u0438\u043c\u0435\u044e\u0449\u0435\u043c\u0441\u044f \u0440\u0430\u0441\u043f\u0440\u0435\u0434\u0435\u043b\u0435\u043d\u0438\u0438, \u043a\u043e\u0442\u043e\u0440\u043e\u0435 \u043d\u0430\u0434\u043e \u0432\u043e\u0441\u043f\u0440\u043e\u0438\u0437\u0432\u0435\u0441\u0442\u0438  *\/ case class OrderBucketsPartitioning(                                      orderExpression: SortOrder,                                      partitionExpressions: Seq[Expression],                                      sortExpressions: Seq[SortOrder],                                      numLines: Int,                                      numPartitions: Int,                                      distribution: Option[Seq[BucketsDistribution]])     extends Expression with Partitioning with Unevaluable {     override def nodeName: String = s\"Repartition with \" +<\/code><\/pre>\n<\/div>\n<\/details>\n<\/div>\n<\/div>\n<\/div>\n<\/div>\n","protected":false},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[],"tags":[],"class_list":["post-381998","post","type-post","status-publish","format-standard","hentry"],"_links":{"self":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/381998","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcomments&post=381998"}],"version-history":[{"count":0,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=\/wp\/v2\/posts\/381998\/revisions"}],"wp:attachment":[{"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fmedia&parent=381998"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Fcategories&post=381998"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/savepearlharbor.com\/index.php?rest_route=%2Fwp%2Fv2%2Ftags&post=381998"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}