Skip to content

Commit 3733726

Browse files
committed
code clean
1 parent c3191a1 commit 3733726

13 files changed

Lines changed: 428 additions & 364 deletions

File tree

core/src/main/java/org/apache/carbondata/core/locks/CarbonLockUtil.java

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,9 @@
1717

1818
package org.apache.carbondata.core.locks;
1919

20+
import java.util.ArrayList;
21+
import java.util.List;
22+
2023
import org.apache.carbondata.common.logging.LogServiceFactory;
2124
import org.apache.carbondata.core.constants.CarbonCommonConstants;
2225
import org.apache.carbondata.core.datastore.filesystem.CarbonFile;
@@ -134,7 +137,6 @@ public static void deleteExpiredSegmentLockFiles(CarbonTable carbonTable) {
134137
}
135138
CarbonFile[] files = FileFactory.getCarbonFile(lockFilesDir)
136139
.listFiles(new CarbonFileFilter() {
137-
138140
@Override
139141
public boolean accept(CarbonFile pathName) {
140142
if (CarbonTablePath.isSegmentLockFilePath(pathName.getName())) {
@@ -149,4 +151,28 @@ public boolean accept(CarbonFile pathName) {
149151
file.delete();
150152
}
151153
}
154+
155+
public static List<ICarbonLock> acquireLock(CarbonTable carbonTable,
156+
List<String> locksToBeAcquired) {
157+
List<ICarbonLock> acquiredLocks = new ArrayList<>();
158+
try {
159+
locksToBeAcquired.forEach(lock ->
160+
acquiredLocks.add(CarbonLockUtil
161+
.getLockObject(carbonTable.getAbsoluteTableIdentifier(), lock))
162+
);
163+
} catch (Exception e) {
164+
releaseLocks(acquiredLocks);
165+
}
166+
return acquiredLocks;
167+
}
168+
169+
public static void releaseLocks(List<ICarbonLock> locks) {
170+
locks.forEach(carbonLock -> {
171+
if (carbonLock.unlock()) {
172+
LOGGER.info("Alter table lock released successfully");
173+
} else {
174+
LOGGER.error("Unable to release lock");
175+
}
176+
});
177+
}
152178
}

core/src/main/java/org/apache/carbondata/core/mutate/CarbonUpdateUtil.java

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -176,7 +176,6 @@ public static String getDeleteDeltaFilePath(String blockPath, String blockName,
176176
String timestamp) {
177177
return blockPath + CarbonCommonConstants.FILE_SEPARATOR + blockName
178178
+ CarbonCommonConstants.HYPHEN + timestamp + CarbonCommonConstants.DELETE_DELTA_FILE_EXT;
179-
180179
}
181180

182181
/**
@@ -280,7 +279,7 @@ public static void mergeSegmentUpdate(boolean isCompaction, List<SegmentUpdateDe
280279
public static boolean updateTableMetadataStatus(Set<Segment> updatedSegmentsList,
281280
CarbonTable table, String updatedTimeStamp, boolean isTimestampUpdateRequired,
282281
boolean isUpdateStatusFileUpdateRequired, List<Segment> segmentsToBeDeleted,
283-
LoadMetadataDetails newLoadEntry) throws IOException{
282+
LoadMetadataDetails newLoadEntry) throws IOException {
284283
return updateTableMetadataStatus(updatedSegmentsList, table, updatedTimeStamp,
285284
isTimestampUpdateRequired, isUpdateStatusFileUpdateRequired,
286285
segmentsToBeDeleted, new ArrayList<Segment>(), "", newLoadEntry);
@@ -298,7 +297,8 @@ public static boolean updateTableMetadataStatus(Set<Segment> updatedSegmentsList
298297
public static boolean updateTableMetadataStatus(Set<Segment> updatedSegmentsList,
299298
CarbonTable table, String updatedTimeStamp, boolean isTimestampUpdateRequired,
300299
boolean isUpdateStatusFileUpdateRequired, List<Segment> segmentsToBeDeleted,
301-
List<Segment> segmentFilesTobeUpdated, String uuid, LoadMetadataDetails newLoadEntry) throws IOException{
300+
List<Segment> segmentFilesTobeUpdated, String uuid,
301+
LoadMetadataDetails newLoadEntry) throws IOException {
302302

303303
boolean status = false;
304304
String metaDataFilepath = table.getMetadataPath();
@@ -370,8 +370,7 @@ public static boolean updateTableMetadataStatus(Set<Segment> updatedSegmentsList
370370
indexToOverwriteNewMetaEntry++;
371371
}
372372
if (!found) {
373-
LOGGER.error("Entry not found to update " + newLoadEntry + " From list :: "
374-
+ listOfLoadFolderDetailsArray);
373+
LOGGER.error("Entry not found to update " + newLoadEntry + " From list");
375374
throw new IOException("Entry not found to update in the table status file");
376375
}
377376
listOfLoadFolderDetailsArray[indexToOverwriteNewMetaEntry] = newLoadEntry;
@@ -814,7 +813,7 @@ public static long readCurrentTime() {
814813
* @param segmentBlockCount
815814
*/
816815
public static void decrementDeletedBlockCount(SegmentUpdateDetails details,
817-
Map<String, Long> segmentBlockCount) {
816+
Map<String, Long> segmentBlockCount) {
818817

819818
String segId = details.getSegmentName();
820819

hadoop/src/main/java/org/apache/carbondata/hadoop/api/CarbonOutputCommitter.java

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
import java.io.IOException;
2121
import java.util.ArrayList;
22-
import java.util.Collections;
2322
import java.util.List;
2423
import java.util.Map;
2524
import java.util.Set;
@@ -35,7 +34,6 @@
3534
import org.apache.carbondata.core.locks.LockUsage;
3635
import org.apache.carbondata.core.metadata.SegmentFileStore;
3736
import org.apache.carbondata.core.metadata.schema.table.CarbonTable;
38-
import org.apache.carbondata.core.mutate.CarbonUpdateUtil;
3937
import org.apache.carbondata.core.statusmanager.LoadMetadataDetails;
4038
import org.apache.carbondata.core.statusmanager.SegmentStatus;
4139
import org.apache.carbondata.core.statusmanager.SegmentStatusManager;
@@ -287,13 +285,16 @@ private void commitJobForPartition(JobContext context, boolean overwriteSet,
287285
newMetaEntry.setDataSize(size);
288286
}
289287
String uniqueId = null;
290-
String updateTime = context.getConfiguration().get(CarbonTableOutputFormat.UPDATE_TIMESTAMP, uniqueId);
288+
String updateTime = context.getConfiguration()
289+
.get(CarbonTableOutputFormat.UPDATE_TIMESTAMP, uniqueId);
291290
if (overwriteSet) {
292291
uniqueId = overwritePartitions(loadModel, newMetaEntry, uuid);
293-
} else if (StringUtils.isNotEmpty(updateTime)){
294-
context.getConfiguration().set("carbon.newMetaEntry", ObjectSerializationUtil.convertObjectToString(newMetaEntry));
292+
} else if (StringUtils.isNotEmpty(updateTime)) {
293+
context.getConfiguration().set("carbon.newMetaEntry",
294+
ObjectSerializationUtil.convertObjectToString(newMetaEntry));
295295
} else {
296-
CarbonLoaderUtil.recordNewLoadMetadata(newMetaEntry, loadModel, false, false, uuid);
296+
CarbonLoaderUtil.recordNewLoadMetadata(newMetaEntry,
297+
loadModel, false, false, uuid);
297298
}
298299
if (operationContext != null) {
299300
operationContext

integration/spark/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAddLoadCommand.scala

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -227,7 +227,6 @@ case class CarbonAddLoadCommand(
227227
model.setDatabaseName(carbonTable.getDatabaseName)
228228
model.setTableName(carbonTable.getTableName)
229229
val operationContext = new OperationContext
230-
operationContext.setProperty("isLoadOrCompaction", false)
231230
val (tableIndexes, indexOperationContext) = CommonLoadUtils.firePreLoadEvents(sparkSession,
232231
model,
233232
"",

integration/spark/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonInsertIntoHadoopFsRelationCommand.scala

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,8 @@ package org.apache.spark.sql.execution.command.management
1919

2020
import java.io.IOException
2121

22-
import org.apache.carbondata.core.statusmanager.LoadMetadataDetails
23-
import org.apache.carbondata.core.util.ObjectSerializationUtil
24-
2522
import scala.collection.mutable
23+
2624
import org.apache.hadoop.fs.{FileSystem, Path}
2725
import org.apache.spark.internal.io.FileCommitProtocol
2826
import org.apache.spark.sql._
@@ -34,9 +32,10 @@ import org.apache.spark.sql.execution.SparkPlan
3432
import org.apache.spark.sql.execution.command._
3533
import org.apache.spark.sql.execution.datasources.{CarbonSQLHadoopMapReduceCommitProtocol, FileFormat, FileFormatWriter, FileIndex, PartitioningUtils, SparkCarbonTableFormat}
3634
import org.apache.spark.sql.internal.SQLConf.PartitionOverwriteMode
35+
import org.apache.spark.sql.types.StringType
3736
import org.apache.spark.sql.util.SchemaUtils
37+
3838
import org.apache.carbondata.spark.util.CarbonScalaUtil
39-
import org.apache.spark.sql.types.StringType
4039

4140
/**
4241
* A command for writing data to a

integration/spark/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonInsertIntoWithDf.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,12 +24,12 @@ import scala.collection.JavaConverters._
2424
import org.apache.spark.sql.{AnalysisException, DataFrame, Row, SparkSession}
2525
import org.apache.spark.sql.execution.command.UpdateTableModel
2626
import org.apache.spark.util.CausedBy
27+
2728
import org.apache.carbondata.common.logging.LogServiceFactory
2829
import org.apache.carbondata.core.datastore.impl.FileFactory
2930
import org.apache.carbondata.core.metadata.schema.table.TableInfo
3031
import org.apache.carbondata.core.statusmanager.{LoadMetadataDetails, SegmentStatus, SegmentStatusManager}
31-
import org.apache.carbondata.core.util.ObjectSerializationUtil.convertStringToObject
32-
import org.apache.carbondata.core.util.{CarbonProperties, ObjectSerializationUtil, ThreadLocalSessionInfo}
32+
import org.apache.carbondata.core.util.{CarbonProperties, ThreadLocalSessionInfo}
3333
import org.apache.carbondata.core.util.path.CarbonTablePath
3434
import org.apache.carbondata.events.OperationContext
3535
import org.apache.carbondata.events.exception.PreEventException

integration/spark/src/main/scala/org/apache/spark/sql/execution/command/management/CommonLoadUtils.scala

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ import java.util
2323
import scala.collection.JavaConverters._
2424
import scala.collection.mutable
2525
import scala.collection.mutable.ArrayBuffer
26+
2627
import org.apache.commons.lang3.StringUtils
2728
import org.apache.hadoop.conf.Configuration
2829
import org.apache.spark.rdd.RDD
@@ -41,6 +42,7 @@ import org.apache.spark.sql.util.SparkSQLUtil
4142
import org.apache.spark.storage.StorageLevel
4243
import org.apache.spark.unsafe.types.UTF8String
4344
import org.apache.spark.util.{CarbonReflectionUtils, CollectionAccumulator, SparkUtil}
45+
4446
import org.apache.carbondata.common.{Maps, Strings}
4547
import org.apache.carbondata.common.logging.LogServiceFactory
4648
import org.apache.carbondata.converter.SparkDataTypeConverterImpl
@@ -1051,11 +1053,11 @@ object CommonLoadUtils {
10511053
ds.collect()
10521054

10531055
if (!loadParams.updateModel.isEmpty && loadParams.updateModel.isDefined
1054-
&& table.isHivePartitionTable) {
1056+
&& table.isHivePartitionTable
1057+
&& !ds.logicalPlan.asInstanceOf[LocalRelation].data(0).anyNull) {
10551058
loadParams.updateModel.get.addedLoadDetail =
10561059
Some(ObjectSerializationUtil.convertStringToObject(
1057-
ds.logicalPlan
1058-
.asInstanceOf[LocalRelation].data(0).getString(0))
1060+
ds.logicalPlan.asInstanceOf[LocalRelation].data(0).getString(0))
10591061
.asInstanceOf[LoadMetadataDetails])
10601062
}
10611063
} catch {

0 commit comments

Comments
 (0)