/
PartitionDropper.scala
123 lines (111 loc) · 5.09 KB
/
PartitionDropper.scala
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.carbondata.spark.rdd
import java.io.IOException
import org.apache.spark.sql.execution.command.{AlterPartitionModel, DropPartitionCallableModel}
import org.apache.spark.util.PartitionUtils
import org.apache.carbondata.common.logging.LogServiceFactory
import org.apache.carbondata.core.metadata.schema.partition.PartitionType
import org.apache.carbondata.spark.{AlterPartitionResultImpl, PartitionFactory}
object PartitionDropper {
val logger = LogServiceFactory.getLogService(PartitionDropper.getClass.getName)
def triggerPartitionDrop(dropPartitionCallableModel: DropPartitionCallableModel): Unit = {
val alterPartitionModel = new AlterPartitionModel(dropPartitionCallableModel.carbonLoadModel,
dropPartitionCallableModel.segmentId.getSegmentNo,
dropPartitionCallableModel.oldPartitionIds,
dropPartitionCallableModel.sqlContext
)
val partitionId = dropPartitionCallableModel.partitionId
val oldPartitionIds = dropPartitionCallableModel.oldPartitionIds
val dropWithData = dropPartitionCallableModel.dropWithData
val carbonTable = dropPartitionCallableModel.carbonTable
val dbName = carbonTable.getDatabaseName
val tableName = carbonTable.getTableName
val absoluteTableIdentifier = carbonTable.getAbsoluteTableIdentifier
val partitionInfo = carbonTable.getPartitionInfo(tableName)
val partitioner = PartitionFactory.getPartitioner(partitionInfo)
var finalDropStatus = false
val bucketInfo = carbonTable.getBucketingInfo(tableName)
val bucketNumber = bucketInfo match {
case null => 1
case _ => bucketInfo.getNumOfRanges
}
val partitionIndex = oldPartitionIds.indexOf(Integer.valueOf(partitionId))
val targetPartitionId = partitionInfo.getPartitionType match {
case PartitionType.RANGE => if (partitionIndex == oldPartitionIds.length - 1) {
"0"
} else {
String.valueOf(oldPartitionIds(partitionIndex + 1))
}
case PartitionType.LIST => "0"
case _ => throw new UnsupportedOperationException(
s"${partitionInfo.getPartitionType} is not supported")
}
if (!dropWithData) {
try {
for (i <- 0 until bucketNumber) {
val bucketId = i
val rdd = new CarbonScanPartitionRDD(alterPartitionModel,
absoluteTableIdentifier,
Seq(partitionId, targetPartitionId),
bucketId
).partitionBy(partitioner).map(_._2)
val dropStatus = new AlterTableLoadPartitionRDD(alterPartitionModel,
new AlterPartitionResultImpl(),
Seq(partitionId),
bucketId,
absoluteTableIdentifier,
rdd).collect()
if (dropStatus.length == 0) {
finalDropStatus = false
} else {
finalDropStatus = dropStatus.forall(_._2)
}
if (!finalDropStatus) {
logger.audit(s"Drop Partition request failed for table " +
s"${ dbName }.${ tableName }")
logger.error(s"Drop Partition request failed for table " +
s"${ dbName }.${ tableName }")
}
}
if (finalDropStatus) {
try {
PartitionUtils.deleteOriginalCarbonFile(alterPartitionModel, absoluteTableIdentifier,
Seq(partitionId, targetPartitionId).toList, dbName,
tableName, partitionInfo)
} catch {
case e: IOException => sys.error(s"Exception while delete original carbon files " +
e.getMessage)
}
logger.audit(s"Drop Partition request completed for table " +
s"${ dbName }.${ tableName }")
logger.info(s"Drop Partition request completed for table " +
s"${ dbName }.${ tableName }")
}
} catch {
case e: Exception => sys.error(s"Exception in dropping partition action: ${ e.getMessage }")
}
} else {
PartitionUtils.deleteOriginalCarbonFile(alterPartitionModel, absoluteTableIdentifier,
Seq(partitionId).toList, dbName, tableName, partitionInfo)
logger.audit(s"Drop Partition request completed for table " +
s"${ dbName }.${ tableName }")
logger.info(s"Drop Partition request completed for table " +
s"${ dbName }.${ tableName }")
}
}
}