You are viewing a plain text version of this content. The canonical link for it is here.
Posted to hdfs-commits@hadoop.apache.org by cu...@apache.org on 2014/08/20 03:34:47 UTC
svn commit: r1619019 [8/11] - in
/hadoop/common/branches/YARN-1051/hadoop-hdfs-project:
hadoop-hdfs-httpfs/src/main/java/org/apache/hadoop/fs/http/client/
hadoop-hdfs-httpfs/src/main/java/org/apache/hadoop/fs/http/server/
hadoop-hdfs-httpfs/src/main/ja...
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeReference.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeReference.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeReference.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeReference.java Wed Aug 20 01:34:29 2014
@@ -287,11 +287,9 @@ public abstract class INodeReference ext
}
@Override
- final INode recordModification(int latestSnapshotId)
+ final void recordModification(int latestSnapshotId)
throws QuotaExceededException {
referred.recordModification(latestSnapshotId);
- // reference is never replaced
- return this;
}
@Override // used by WithCount
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeSymlink.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeSymlink.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeSymlink.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/INodeSymlink.java Wed Aug 20 01:34:29 2014
@@ -47,12 +47,11 @@ public class INodeSymlink extends INodeW
}
@Override
- INode recordModification(int latestSnapshotId) throws QuotaExceededException {
+ void recordModification(int latestSnapshotId) throws QuotaExceededException {
if (isInLatestSnapshot(latestSnapshotId)) {
INodeDirectory parent = getParent();
parent.saveChild2Snapshot(this, latestSnapshotId, new INodeSymlink(this));
}
- return this;
}
/** @return true unconditionally. */
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/NameNodeRpcServer.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/NameNodeRpcServer.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/NameNodeRpcServer.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/NameNodeRpcServer.java Wed Aug 20 01:34:29 2014
@@ -54,6 +54,7 @@ import org.apache.hadoop.fs.XAttrSetFlag
import org.apache.hadoop.fs.permission.AclEntry;
import org.apache.hadoop.fs.permission.AclStatus;
import org.apache.hadoop.fs.permission.FsPermission;
+import org.apache.hadoop.fs.permission.FsAction;
import org.apache.hadoop.fs.permission.PermissionStatus;
import org.apache.hadoop.ha.HAServiceStatus;
import org.apache.hadoop.ha.HealthCheckFailedException;
@@ -115,6 +116,7 @@ import org.apache.hadoop.hdfs.server.pro
import org.apache.hadoop.hdfs.server.protocol.DatanodeCommand;
import org.apache.hadoop.hdfs.server.protocol.DatanodeProtocol;
import org.apache.hadoop.hdfs.server.protocol.DatanodeRegistration;
+import org.apache.hadoop.hdfs.server.protocol.DatanodeStorageReport;
import org.apache.hadoop.hdfs.server.protocol.FinalizeCommand;
import org.apache.hadoop.hdfs.server.protocol.HeartbeatResponse;
import org.apache.hadoop.hdfs.server.protocol.NamenodeCommand;
@@ -830,12 +832,24 @@ class NameNodeRpcServer implements Namen
throws IOException {
DatanodeInfo results[] = namesystem.datanodeReport(type);
if (results == null ) {
- throw new IOException("Cannot find datanode report");
+ throw new IOException("Failed to get datanode report for " + type
+ + " datanodes.");
}
return results;
}
@Override // ClientProtocol
+ public DatanodeStorageReport[] getDatanodeStorageReport(
+ DatanodeReportType type) throws IOException {
+ final DatanodeStorageReport[] reports = namesystem.getDatanodeStorageReport(type);
+ if (reports == null ) {
+ throw new IOException("Failed to get datanode storage report for " + type
+ + " datanodes.");
+ }
+ return reports;
+ }
+
+ @Override // ClientProtocol
public boolean setSafeMode(SafeModeAction action, boolean isChecked)
throws IOException {
OperationCategory opCategory = OperationCategory.UNCHECKED;
@@ -1051,7 +1065,7 @@ class NameNodeRpcServer implements Namen
// for the same node and storage, so the value returned by the last
// call of this loop is the final updated value for noStaleStorage.
//
- noStaleStorages = bm.processReport(nodeReg, r.getStorage(), poolId, blocks);
+ noStaleStorages = bm.processReport(nodeReg, r.getStorage(), blocks);
metrics.incrStorageBlockReportOps();
}
@@ -1087,7 +1101,7 @@ class NameNodeRpcServer implements Namen
+" blocks.");
}
for(StorageReceivedDeletedBlocks r : receivedAndDeletedBlocks) {
- namesystem.processIncrementalBlockReport(nodeReg, poolId, r);
+ namesystem.processIncrementalBlockReport(nodeReg, r);
}
}
@@ -1430,5 +1444,10 @@ class NameNodeRpcServer implements Namen
public void removeXAttr(String src, XAttr xAttr) throws IOException {
namesystem.removeXAttr(src, xAttr);
}
+
+ @Override
+ public void checkAccess(String path, FsAction mode) throws IOException {
+ namesystem.checkAccess(path, mode);
+ }
}
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/metrics/FSNamesystemMBean.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/metrics/FSNamesystemMBean.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/metrics/FSNamesystemMBean.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/metrics/FSNamesystemMBean.java Wed Aug 20 01:34:29 2014
@@ -151,4 +151,11 @@ public interface FSNamesystemMBean {
* @return number of blocks pending deletion
*/
long getPendingDeletionBlocks();
+
+ /**
+ * Number of content stale storages.
+ * @return number of content stale storages
+ */
+ public int getNumStaleStorages();
+
}
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/web/resources/NamenodeWebHdfsMethods.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/web/resources/NamenodeWebHdfsMethods.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/web/resources/NamenodeWebHdfsMethods.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/namenode/web/resources/NamenodeWebHdfsMethods.java Wed Aug 20 01:34:29 2014
@@ -57,6 +57,7 @@ import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.Options;
import org.apache.hadoop.fs.XAttr;
import org.apache.hadoop.fs.permission.AclStatus;
+import org.apache.hadoop.fs.permission.FsAction;
import org.apache.hadoop.hdfs.StorageType;
import org.apache.hadoop.hdfs.XAttrHelper;
import org.apache.hadoop.hdfs.protocol.DatanodeInfo;
@@ -112,6 +113,7 @@ import org.apache.hadoop.hdfs.web.resour
import org.apache.hadoop.hdfs.web.resources.XAttrNameParam;
import org.apache.hadoop.hdfs.web.resources.XAttrSetFlagParam;
import org.apache.hadoop.hdfs.web.resources.XAttrValueParam;
+import org.apache.hadoop.hdfs.web.resources.FsActionParam;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.ipc.RetriableException;
import org.apache.hadoop.ipc.Server;
@@ -755,10 +757,12 @@ public class NamenodeWebHdfsMethods {
@QueryParam(XAttrEncodingParam.NAME) @DefaultValue(XAttrEncodingParam.DEFAULT)
final XAttrEncodingParam xattrEncoding,
@QueryParam(ExcludeDatanodesParam.NAME) @DefaultValue(ExcludeDatanodesParam.DEFAULT)
- final ExcludeDatanodesParam excludeDatanodes
+ final ExcludeDatanodesParam excludeDatanodes,
+ @QueryParam(FsActionParam.NAME) @DefaultValue(FsActionParam.DEFAULT)
+ final FsActionParam fsAction
) throws IOException, InterruptedException {
return get(ugi, delegation, username, doAsUser, ROOT, op, offset, length,
- renewer, bufferSize, xattrNames, xattrEncoding, excludeDatanodes);
+ renewer, bufferSize, xattrNames, xattrEncoding, excludeDatanodes, fsAction);
}
/** Handle HTTP GET request. */
@@ -789,11 +793,13 @@ public class NamenodeWebHdfsMethods {
@QueryParam(XAttrEncodingParam.NAME) @DefaultValue(XAttrEncodingParam.DEFAULT)
final XAttrEncodingParam xattrEncoding,
@QueryParam(ExcludeDatanodesParam.NAME) @DefaultValue(ExcludeDatanodesParam.DEFAULT)
- final ExcludeDatanodesParam excludeDatanodes
+ final ExcludeDatanodesParam excludeDatanodes,
+ @QueryParam(FsActionParam.NAME) @DefaultValue(FsActionParam.DEFAULT)
+ final FsActionParam fsAction
) throws IOException, InterruptedException {
init(ugi, delegation, username, doAsUser, path, op, offset, length,
- renewer, bufferSize, xattrEncoding, excludeDatanodes);
+ renewer, bufferSize, xattrEncoding, excludeDatanodes, fsAction);
return ugi.doAs(new PrivilegedExceptionAction<Response>() {
@Override
@@ -801,7 +807,7 @@ public class NamenodeWebHdfsMethods {
try {
return get(ugi, delegation, username, doAsUser,
path.getAbsolutePath(), op, offset, length, renewer, bufferSize,
- xattrNames, xattrEncoding, excludeDatanodes);
+ xattrNames, xattrEncoding, excludeDatanodes, fsAction);
} finally {
reset();
}
@@ -822,7 +828,8 @@ public class NamenodeWebHdfsMethods {
final BufferSizeParam bufferSize,
final List<XAttrNameParam> xattrNames,
final XAttrEncodingParam xattrEncoding,
- final ExcludeDatanodesParam excludeDatanodes
+ final ExcludeDatanodesParam excludeDatanodes,
+ final FsActionParam fsAction
) throws IOException, URISyntaxException {
final NameNode namenode = (NameNode)context.getAttribute("name.node");
final NamenodeProtocols np = getRPCServer(namenode);
@@ -919,6 +926,10 @@ public class NamenodeWebHdfsMethods {
final String js = JsonUtil.toJsonString(xAttrs);
return Response.ok(js).type(MediaType.APPLICATION_JSON).build();
}
+ case CHECKACCESS: {
+ np.checkAccess(fullpath, FsAction.getFsAction(fsAction.getValue()));
+ return Response.ok().build();
+ }
default:
throw new UnsupportedOperationException(op + " is not supported");
}
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/BlocksWithLocations.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/BlocksWithLocations.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/BlocksWithLocations.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/BlocksWithLocations.java Wed Aug 20 01:34:29 2014
@@ -17,10 +17,9 @@
*/
package org.apache.hadoop.hdfs.server.protocol;
-import java.util.Arrays;
-
import org.apache.hadoop.classification.InterfaceAudience;
import org.apache.hadoop.classification.InterfaceStability;
+import org.apache.hadoop.hdfs.StorageType;
import org.apache.hadoop.hdfs.protocol.Block;
/**
@@ -39,12 +38,15 @@ public class BlocksWithLocations {
final Block block;
final String[] datanodeUuids;
final String[] storageIDs;
+ final StorageType[] storageTypes;
/** constructor */
- public BlockWithLocations(Block block, String[] datanodeUuids, String[] storageIDs) {
+ public BlockWithLocations(Block block, String[] datanodeUuids,
+ String[] storageIDs, StorageType[] storageTypes) {
this.block = block;
this.datanodeUuids = datanodeUuids;
this.storageIDs = storageIDs;
+ this.storageTypes = storageTypes;
}
/** get the block */
@@ -61,7 +63,12 @@ public class BlocksWithLocations {
public String[] getStorageIDs() {
return storageIDs;
}
-
+
+ /** @return the storage types */
+ public StorageType[] getStorageTypes() {
+ return storageTypes;
+ }
+
@Override
public String toString() {
final StringBuilder b = new StringBuilder();
@@ -70,12 +77,18 @@ public class BlocksWithLocations {
return b.append("[]").toString();
}
- b.append(storageIDs[0]).append('@').append(datanodeUuids[0]);
+ appendString(0, b.append("["));
for(int i = 1; i < datanodeUuids.length; i++) {
- b.append(", ").append(storageIDs[i]).append("@").append(datanodeUuids[i]);
+ appendString(i, b.append(","));
}
return b.append("]").toString();
}
+
+ private StringBuilder appendString(int i, StringBuilder b) {
+ return b.append("[").append(storageTypes[i]).append("]")
+ .append(storageIDs[i])
+ .append("@").append(datanodeUuids[i]);
+ }
}
private final BlockWithLocations[] blocks;
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/RegisterCommand.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/RegisterCommand.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/RegisterCommand.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/protocol/RegisterCommand.java Wed Aug 20 01:34:29 2014
@@ -22,6 +22,9 @@ import org.apache.hadoop.classification.
/**
* A BlockCommand is an instruction to a datanode to register with the namenode.
+ * This command can't be combined with other commands in the same response.
+ * This is because after the datanode processes RegisterCommand, it will skip
+ * the rest of the DatanodeCommands in the same HeartbeatResponse.
*/
@InterfaceAudience.Private
@InterfaceStability.Evolving
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShm.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShm.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShm.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShm.java Wed Aug 20 01:34:29 2014
@@ -32,11 +32,16 @@ import com.google.common.base.Preconditi
* DfsClientShm is a subclass of ShortCircuitShm which is used by the
* DfsClient.
* When the UNIX domain socket associated with this shared memory segment
- * closes unexpectedly, we mark the slots inside this segment as stale.
- * ShortCircuitReplica objects that contain stale slots are themselves stale,
+ * closes unexpectedly, we mark the slots inside this segment as disconnected.
+ * ShortCircuitReplica objects that contain disconnected slots are stale,
* and will not be used to service new reads or mmap operations.
* However, in-progress read or mmap operations will continue to proceed.
* Once the last slot is deallocated, the segment can be safely munmapped.
+ *
+ * Slots may also become stale because the associated replica has been deleted
+ * on the DataNode. In this case, the DataNode will clear the 'valid' bit.
+ * The client will then see these slots as stale (see
+ * #{ShortCircuitReplica#isStale}).
*/
public class DfsClientShm extends ShortCircuitShm
implements DomainSocketWatcher.Handler {
@@ -58,7 +63,7 @@ public class DfsClientShm extends ShortC
*
* {@link DfsClientShm#handle} sets this to true.
*/
- private boolean stale = false;
+ private boolean disconnected = false;
DfsClientShm(ShmId shmId, FileInputStream stream, EndpointShmManager manager,
DomainPeer peer) throws IOException {
@@ -76,14 +81,14 @@ public class DfsClientShm extends ShortC
}
/**
- * Determine if the shared memory segment is stale.
+ * Determine if the shared memory segment is disconnected from the DataNode.
*
* This must be called with the DfsClientShmManager lock held.
*
* @return True if the shared memory segment is stale.
*/
- public synchronized boolean isStale() {
- return stale;
+ public synchronized boolean isDisconnected() {
+ return disconnected;
}
/**
@@ -97,8 +102,8 @@ public class DfsClientShm extends ShortC
public boolean handle(DomainSocket sock) {
manager.unregisterShm(getShmId());
synchronized (this) {
- Preconditions.checkState(!stale);
- stale = true;
+ Preconditions.checkState(!disconnected);
+ disconnected = true;
boolean hadSlots = false;
for (Iterator<Slot> iter = slotIterator(); iter.hasNext(); ) {
Slot slot = iter.next();
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShmManager.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShmManager.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShmManager.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/DfsClientShmManager.java Wed Aug 20 01:34:29 2014
@@ -271,12 +271,12 @@ public class DfsClientShmManager impleme
loading = false;
finishedLoading.signalAll();
}
- if (shm.isStale()) {
+ if (shm.isDisconnected()) {
// If the peer closed immediately after the shared memory segment
// was created, the DomainSocketWatcher callback might already have
- // fired and marked the shm as stale. In this case, we obviously
- // don't want to add the SharedMemorySegment to our list of valid
- // not-full segments.
+ // fired and marked the shm as disconnected. In this case, we
+ // obviously don't want to add the SharedMemorySegment to our list
+ // of valid not-full segments.
if (LOG.isDebugEnabled()) {
LOG.debug(this + ": the UNIX domain socket associated with " +
"this short-circuit memory closed before we could make " +
@@ -299,7 +299,7 @@ public class DfsClientShmManager impleme
void freeSlot(Slot slot) {
DfsClientShm shm = (DfsClientShm)slot.getShm();
shm.unregisterSlot(slot.getSlotIdx());
- if (shm.isStale()) {
+ if (shm.isDisconnected()) {
// Stale shared memory segments should not be tracked here.
Preconditions.checkState(!full.containsKey(shm.getShmId()));
Preconditions.checkState(!notFull.containsKey(shm.getShmId()));
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/ShortCircuitShm.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/ShortCircuitShm.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/ShortCircuitShm.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/shortcircuit/ShortCircuitShm.java Wed Aug 20 01:34:29 2014
@@ -306,6 +306,13 @@ public class ShortCircuitShm {
(slotAddress - baseAddress) / BYTES_PER_SLOT);
}
+ /**
+ * Clear the slot.
+ */
+ void clear() {
+ unsafe.putLongVolatile(null, this.slotAddress, 0);
+ }
+
private boolean isSet(long flag) {
long prev = unsafe.getLongVolatile(null, this.slotAddress);
return (prev & flag) != 0;
@@ -535,6 +542,7 @@ public class ShortCircuitShm {
}
allocatedSlots.set(idx, true);
Slot slot = new Slot(calculateSlotAddress(idx), blockId);
+ slot.clear();
slot.makeValid();
slots[idx] = slot;
if (LOG.isTraceEnabled()) {
@@ -583,7 +591,7 @@ public class ShortCircuitShm {
Slot slot = new Slot(calculateSlotAddress(slotIdx), blockId);
if (!slot.isValid()) {
throw new InvalidRequestException(this + ": slot " + slotIdx +
- " has not been allocated.");
+ " is not marked as valid.");
}
slots[slotIdx] = slot;
allocatedSlots.set(slotIdx, true);
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/tools/offlineEditsViewer/XmlEditsVisitor.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/tools/offlineEditsViewer/XmlEditsVisitor.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/tools/offlineEditsViewer/XmlEditsVisitor.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/tools/offlineEditsViewer/XmlEditsVisitor.java Wed Aug 20 01:34:29 2014
@@ -29,8 +29,8 @@ import org.xml.sax.ContentHandler;
import org.xml.sax.SAXException;
import org.xml.sax.helpers.AttributesImpl;
-import com.sun.org.apache.xml.internal.serialize.OutputFormat;
-import com.sun.org.apache.xml.internal.serialize.XMLSerializer;
+import org.apache.xml.serialize.OutputFormat;
+import org.apache.xml.serialize.XMLSerializer;
/**
* An XmlEditsVisitor walks over an EditLog structure and writes out
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/util/EnumCounters.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/util/EnumCounters.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/util/EnumCounters.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/util/EnumCounters.java Wed Aug 20 01:34:29 2014
@@ -37,7 +37,7 @@ import com.google.common.base.Preconditi
public class EnumCounters<E extends Enum<E>> {
/** The class of the enum. */
private final Class<E> enumClass;
- /** The counter array, counters[i] corresponds to the enumConstants[i]. */
+ /** An array of longs corresponding to the enum type. */
private final long[] counters;
/**
@@ -75,6 +75,13 @@ public class EnumCounters<E extends Enum
}
}
+ /** Reset all counters to zero. */
+ public final void reset() {
+ for(int i = 0; i < counters.length; i++) {
+ this.counters[i] = 0L;
+ }
+ }
+
/** Add the given value to counter e. */
public final void add(final E e, final long value) {
counters[e.ordinal()] += value;
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/WebHdfsFileSystem.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/WebHdfsFileSystem.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/WebHdfsFileSystem.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/WebHdfsFileSystem.java Wed Aug 20 01:34:29 2014
@@ -54,6 +54,7 @@ import org.apache.hadoop.fs.XAttrCodec;
import org.apache.hadoop.fs.XAttrSetFlag;
import org.apache.hadoop.fs.permission.AclEntry;
import org.apache.hadoop.fs.permission.AclStatus;
+import org.apache.hadoop.fs.permission.FsAction;
import org.apache.hadoop.fs.permission.FsPermission;
import org.apache.hadoop.hdfs.DFSConfigKeys;
import org.apache.hadoop.hdfs.DFSUtil;
@@ -1357,6 +1358,12 @@ public class WebHdfsFileSystem extends F
}
@Override
+ public void access(final Path path, final FsAction mode) throws IOException {
+ final HttpOpParam.Op op = GetOpParam.Op.CHECKACCESS;
+ new FsPathRunner(op, path, new FsActionParam(mode)).run();
+ }
+
+ @Override
public ContentSummary getContentSummary(final Path p) throws IOException {
statistics.incrementReadOps(1);
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/resources/GetOpParam.java
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/resources/GetOpParam.java?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/resources/GetOpParam.java (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/web/resources/GetOpParam.java Wed Aug 20 01:34:29 2014
@@ -39,7 +39,9 @@ public class GetOpParam extends HttpOpPa
GETXATTRS(false, HttpURLConnection.HTTP_OK),
LISTXATTRS(false, HttpURLConnection.HTTP_OK),
- NULL(false, HttpURLConnection.HTTP_NOT_IMPLEMENTED);
+ NULL(false, HttpURLConnection.HTTP_NOT_IMPLEMENTED),
+
+ CHECKACCESS(false, HttpURLConnection.HTTP_OK);
final boolean redirect;
final int expectedHttpResponseCode;
Propchange: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/
------------------------------------------------------------------------------
Merged /hadoop/common/trunk/hadoop-hdfs-project/hadoop-hdfs/src/main/native:r1603348-1619017
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/fuse-dfs/CMakeLists.txt
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/fuse-dfs/CMakeLists.txt?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/fuse-dfs/CMakeLists.txt (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/fuse-dfs/CMakeLists.txt Wed Aug 20 01:34:29 2014
@@ -37,6 +37,10 @@ ELSE (${CMAKE_SYSTEM_NAME} MATCHES "Linu
ENDIF (${CMAKE_SYSTEM_NAME} MATCHES "Linux")
IF(FUSE_FOUND)
+ add_library(posix_util
+ ../util/posix_util.c
+ )
+
add_executable(fuse_dfs
fuse_dfs.c
fuse_options.c
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.c
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.c?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.c (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.c Wed Aug 20 01:34:29 2014
@@ -19,8 +19,8 @@
#include "exception.h"
#include "hdfs.h"
#include "jni_helper.h"
+#include "platform.h"
-#include <inttypes.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
@@ -35,54 +35,54 @@ struct ExceptionInfo {
static const struct ExceptionInfo gExceptionInfo[] = {
{
- .name = "java.io.FileNotFoundException",
- .noPrintFlag = NOPRINT_EXC_FILE_NOT_FOUND,
- .excErrno = ENOENT,
+ "java.io.FileNotFoundException",
+ NOPRINT_EXC_FILE_NOT_FOUND,
+ ENOENT,
},
{
- .name = "org.apache.hadoop.security.AccessControlException",
- .noPrintFlag = NOPRINT_EXC_ACCESS_CONTROL,
- .excErrno = EACCES,
+ "org.apache.hadoop.security.AccessControlException",
+ NOPRINT_EXC_ACCESS_CONTROL,
+ EACCES,
},
{
- .name = "org.apache.hadoop.fs.UnresolvedLinkException",
- .noPrintFlag = NOPRINT_EXC_UNRESOLVED_LINK,
- .excErrno = ENOLINK,
+ "org.apache.hadoop.fs.UnresolvedLinkException",
+ NOPRINT_EXC_UNRESOLVED_LINK,
+ ENOLINK,
},
{
- .name = "org.apache.hadoop.fs.ParentNotDirectoryException",
- .noPrintFlag = NOPRINT_EXC_PARENT_NOT_DIRECTORY,
- .excErrno = ENOTDIR,
+ "org.apache.hadoop.fs.ParentNotDirectoryException",
+ NOPRINT_EXC_PARENT_NOT_DIRECTORY,
+ ENOTDIR,
},
{
- .name = "java.lang.IllegalArgumentException",
- .noPrintFlag = NOPRINT_EXC_ILLEGAL_ARGUMENT,
- .excErrno = EINVAL,
+ "java.lang.IllegalArgumentException",
+ NOPRINT_EXC_ILLEGAL_ARGUMENT,
+ EINVAL,
},
{
- .name = "java.lang.OutOfMemoryError",
- .noPrintFlag = 0,
- .excErrno = ENOMEM,
+ "java.lang.OutOfMemoryError",
+ 0,
+ ENOMEM,
},
{
- .name = "org.apache.hadoop.hdfs.server.namenode.SafeModeException",
- .noPrintFlag = 0,
- .excErrno = EROFS,
+ "org.apache.hadoop.hdfs.server.namenode.SafeModeException",
+ 0,
+ EROFS,
},
{
- .name = "org.apache.hadoop.fs.FileAlreadyExistsException",
- .noPrintFlag = 0,
- .excErrno = EEXIST,
+ "org.apache.hadoop.fs.FileAlreadyExistsException",
+ 0,
+ EEXIST,
},
{
- .name = "org.apache.hadoop.hdfs.protocol.QuotaExceededException",
- .noPrintFlag = 0,
- .excErrno = EDQUOT,
+ "org.apache.hadoop.hdfs.protocol.QuotaExceededException",
+ 0,
+ EDQUOT,
},
{
- .name = "org.apache.hadoop.hdfs.server.namenode.LeaseExpiredException",
- .noPrintFlag = 0,
- .excErrno = ESTALE,
+ "org.apache.hadoop.hdfs.server.namenode.LeaseExpiredException",
+ 0,
+ ESTALE,
},
};
@@ -113,6 +113,7 @@ int printExceptionAndFreeV(JNIEnv *env,
jstring jStr = NULL;
jvalue jVal;
jthrowable jthr;
+ const char *stackTrace;
jthr = classNameOfObject(exc, env, &className);
if (jthr) {
@@ -148,7 +149,7 @@ int printExceptionAndFreeV(JNIEnv *env,
destroyLocalReference(env, jthr);
} else {
jStr = jVal.l;
- const char *stackTrace = (*env)->GetStringUTFChars(env, jStr, NULL);
+ stackTrace = (*env)->GetStringUTFChars(env, jStr, NULL);
if (!stackTrace) {
fprintf(stderr, "(unable to get stack trace for %s exception: "
"GetStringUTFChars error.)\n", className);
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.h
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.h?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.h (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/exception.h Wed Aug 20 01:34:29 2014
@@ -34,13 +34,14 @@
* usually not what you want.)
*/
+#include "platform.h"
+
#include <jni.h>
#include <stdio.h>
#include <stdlib.h>
#include <stdarg.h>
#include <search.h>
-#include <pthread.h>
#include <errno.h>
/**
@@ -109,7 +110,7 @@ int printExceptionAndFreeV(JNIEnv *env,
* object.
*/
int printExceptionAndFree(JNIEnv *env, jthrowable exc, int noPrintFlags,
- const char *fmt, ...) __attribute__((format(printf, 4, 5)));
+ const char *fmt, ...) TYPE_CHECKED_PRINTF_FORMAT(4, 5);
/**
* Print out information about the pending exception and free it.
@@ -124,7 +125,7 @@ int printExceptionAndFree(JNIEnv *env, j
* object.
*/
int printPendingExceptionAndFree(JNIEnv *env, int noPrintFlags,
- const char *fmt, ...) __attribute__((format(printf, 3, 4)));
+ const char *fmt, ...) TYPE_CHECKED_PRINTF_FORMAT(3, 4);
/**
* Get a local reference to the pending exception and clear it.
@@ -150,6 +151,7 @@ jthrowable getPendingExceptionAndClear(J
* @return A local reference to a RuntimeError
*/
jthrowable newRuntimeError(JNIEnv *env, const char *fmt, ...)
- __attribute__((format(printf, 2, 3)));
+ TYPE_CHECKED_PRINTF_FORMAT(2, 3);
+#undef TYPE_CHECKED_PRINTF_FORMAT
#endif
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.c
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.c?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.c (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.c Wed Aug 20 01:34:29 2014
@@ -49,18 +49,18 @@ int expectFileStats(hdfsFile file,
stats->totalShortCircuitBytesRead,
stats->totalZeroCopyBytesRead);
if (expectedTotalBytesRead != UINT64_MAX) {
- EXPECT_INT64_EQ(expectedTotalBytesRead, stats->totalBytesRead);
+ EXPECT_UINT64_EQ(expectedTotalBytesRead, stats->totalBytesRead);
}
if (expectedTotalLocalBytesRead != UINT64_MAX) {
- EXPECT_INT64_EQ(expectedTotalLocalBytesRead,
+ EXPECT_UINT64_EQ(expectedTotalLocalBytesRead,
stats->totalLocalBytesRead);
}
if (expectedTotalShortCircuitBytesRead != UINT64_MAX) {
- EXPECT_INT64_EQ(expectedTotalShortCircuitBytesRead,
+ EXPECT_UINT64_EQ(expectedTotalShortCircuitBytesRead,
stats->totalShortCircuitBytesRead);
}
if (expectedTotalZeroCopyBytesRead != UINT64_MAX) {
- EXPECT_INT64_EQ(expectedTotalZeroCopyBytesRead,
+ EXPECT_UINT64_EQ(expectedTotalZeroCopyBytesRead,
stats->totalZeroCopyBytesRead);
}
hdfsFileFreeReadStatistics(stats);
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.h
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.h?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.h (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/expect.h Wed Aug 20 01:34:29 2014
@@ -126,6 +126,18 @@ struct hdfsFile_internal;
} \
} while (0);
+#define EXPECT_UINT64_EQ(x, y) \
+ do { \
+ uint64_t __my_ret__ = y; \
+ int __my_errno__ = errno; \
+ if (__my_ret__ != (x)) { \
+ fprintf(stderr, "TEST_ERROR: failed on %s:%d with return " \
+ "value %"PRIu64" (errno: %d): expected %"PRIu64"\n", \
+ __FILE__, __LINE__, __my_ret__, __my_errno__, (x)); \
+ return -1; \
+ } \
+ } while (0);
+
#define RETRY_ON_EINTR_GET_ERRNO(ret, expr) do { \
ret = expr; \
if (!ret) \
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/hdfs.c
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/hdfs.c?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/hdfs.c (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/hdfs.c Wed Aug 20 01:34:29 2014
@@ -19,7 +19,9 @@
#include "exception.h"
#include "hdfs.h"
#include "jni_helper.h"
+#include "platform.h"
+#include <fcntl.h>
#include <inttypes.h>
#include <stdio.h>
#include <string.h>
@@ -63,9 +65,9 @@ static void hdfsFreeFileInfoEntry(hdfsFi
*/
enum hdfsStreamType
{
- UNINITIALIZED = 0,
- INPUT = 1,
- OUTPUT = 2,
+ HDFS_STREAM_UNINITIALIZED = 0,
+ HDFS_STREAM_INPUT = 1,
+ HDFS_STREAM_OUTPUT = 2,
};
/**
@@ -79,7 +81,7 @@ struct hdfsFile_internal {
int hdfsFileIsOpenForRead(hdfsFile file)
{
- return (file->type == INPUT);
+ return (file->type == HDFS_STREAM_INPUT);
}
int hdfsFileGetReadStatistics(hdfsFile file,
@@ -96,7 +98,7 @@ int hdfsFileGetReadStatistics(hdfsFile f
errno = EINTERNAL;
return -1;
}
- if (file->type != INPUT) {
+ if (file->type != HDFS_STREAM_INPUT) {
ret = EINVAL;
goto done;
}
@@ -180,7 +182,7 @@ void hdfsFileFreeReadStatistics(struct h
int hdfsFileIsOpenForWrite(hdfsFile file)
{
- return (file->type == OUTPUT);
+ return (file->type == HDFS_STREAM_OUTPUT);
}
int hdfsFileUsesDirectRead(hdfsFile file)
@@ -441,7 +443,7 @@ void hdfsBuilderSetKerbTicketCachePath(s
bld->kerbTicketCachePath = kerbTicketCachePath;
}
-hdfsFS hdfsConnect(const char* host, tPort port)
+hdfsFS hdfsConnect(const char *host, tPort port)
{
struct hdfsBuilder *bld = hdfsNewBuilder();
if (!bld)
@@ -452,7 +454,7 @@ hdfsFS hdfsConnect(const char* host, tPo
}
/** Always return a new FileSystem handle */
-hdfsFS hdfsConnectNewInstance(const char* host, tPort port)
+hdfsFS hdfsConnectNewInstance(const char *host, tPort port)
{
struct hdfsBuilder *bld = hdfsNewBuilder();
if (!bld)
@@ -463,7 +465,7 @@ hdfsFS hdfsConnectNewInstance(const char
return hdfsBuilderConnect(bld);
}
-hdfsFS hdfsConnectAsUser(const char* host, tPort port, const char *user)
+hdfsFS hdfsConnectAsUser(const char *host, tPort port, const char *user)
{
struct hdfsBuilder *bld = hdfsNewBuilder();
if (!bld)
@@ -475,7 +477,7 @@ hdfsFS hdfsConnectAsUser(const char* hos
}
/** Always return a new FileSystem handle */
-hdfsFS hdfsConnectAsUserNewInstance(const char* host, tPort port,
+hdfsFS hdfsConnectAsUserNewInstance(const char *host, tPort port,
const char *user)
{
struct hdfsBuilder *bld = hdfsNewBuilder();
@@ -518,7 +520,7 @@ static int calcEffectiveURI(struct hdfsB
if (bld->port == 0) {
suffix[0] = '\0';
} else {
- lastColon = rindex(bld->nn, ':');
+ lastColon = strrchr(bld->nn, ':');
if (lastColon && (strspn(lastColon + 1, "0123456789") ==
strlen(lastColon + 1))) {
fprintf(stderr, "port %d was given, but URI '%s' already "
@@ -737,6 +739,8 @@ int hdfsDisconnect(hdfsFS fs)
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
int ret;
+ jobject jFS;
+ jthrowable jthr;
if (env == NULL) {
errno = EINTERNAL;
@@ -744,7 +748,7 @@ int hdfsDisconnect(hdfsFS fs)
}
//Parameters
- jobject jFS = (jobject)fs;
+ jFS = (jobject)fs;
//Sanity check
if (fs == NULL) {
@@ -752,7 +756,7 @@ int hdfsDisconnect(hdfsFS fs)
return -1;
}
- jthrowable jthr = invokeMethod(env, NULL, INSTANCE, jFS, HADOOP_FS,
+ jthr = invokeMethod(env, NULL, INSTANCE, jFS, HADOOP_FS,
"close", "()V");
if (jthr) {
ret = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -792,7 +796,7 @@ static jthrowable getDefaultBlockSize(JN
return NULL;
}
-hdfsFile hdfsOpenFile(hdfsFS fs, const char* path, int flags,
+hdfsFile hdfsOpenFile(hdfsFS fs, const char *path, int flags,
int bufferSize, short replication, tSize blockSize)
{
/*
@@ -801,15 +805,7 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
FSData{Input|Output}Stream f{is|os} = fs.create(f);
return f{is|os};
*/
- /* Get the JNIEnv* corresponding to current thread */
- JNIEnv* env = getJNIEnv();
int accmode = flags & O_ACCMODE;
-
- if (env == NULL) {
- errno = EINTERNAL;
- return NULL;
- }
-
jstring jStrBufferSize = NULL, jStrReplication = NULL;
jobject jConfiguration = NULL, jPath = NULL, jFile = NULL;
jobject jFS = (jobject)fs;
@@ -817,6 +813,20 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
jvalue jVal;
hdfsFile file = NULL;
int ret;
+ jint jBufferSize = bufferSize;
+ jshort jReplication = replication;
+
+ /* The hadoop java api/signature */
+ const char *method = NULL;
+ const char *signature = NULL;
+
+ /* Get the JNIEnv* corresponding to current thread */
+ JNIEnv* env = getJNIEnv();
+ if (env == NULL) {
+ errno = EINTERNAL;
+ return NULL;
+ }
+
if (accmode == O_RDONLY || accmode == O_WRONLY) {
/* yay */
@@ -834,10 +844,6 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
fprintf(stderr, "WARN: hdfs does not truly support O_CREATE && O_EXCL\n");
}
- /* The hadoop java api/signature */
- const char* method = NULL;
- const char* signature = NULL;
-
if (accmode == O_RDONLY) {
method = "open";
signature = JMETHOD2(JPARAM(HADOOP_PATH), "I", JPARAM(HADOOP_ISTRM));
@@ -867,8 +873,6 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
}
jConfiguration = jVal.l;
- jint jBufferSize = bufferSize;
- jshort jReplication = replication;
jStrBufferSize = (*env)->NewStringUTF(env, "io.file.buffer.size");
if (!jStrBufferSize) {
ret = printPendingExceptionAndFree(env, PRINT_EXC_ALL, "OOM");
@@ -905,7 +909,7 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
path);
goto done;
}
- jReplication = jVal.i;
+ jReplication = (jshort)jVal.i;
}
}
@@ -955,7 +959,8 @@ hdfsFile hdfsOpenFile(hdfsFS fs, const c
"hdfsOpenFile(%s): NewGlobalRef", path);
goto done;
}
- file->type = (((flags & O_WRONLY) == 0) ? INPUT : OUTPUT);
+ file->type = (((flags & O_WRONLY) == 0) ? HDFS_STREAM_INPUT :
+ HDFS_STREAM_OUTPUT);
file->flags = 0;
if ((flags & O_WRONLY) == 0) {
@@ -998,31 +1003,33 @@ int hdfsCloseFile(hdfsFS fs, hdfsFile fi
// JAVA EQUIVALENT:
// file.close
+ //The interface whose 'close' method to be called
+ const char *interface;
+ const char *interfaceShortName;
+
+ //Caught exception
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
-
if (env == NULL) {
errno = EINTERNAL;
return -1;
}
- //Caught exception
- jthrowable jthr;
-
//Sanity check
- if (!file || file->type == UNINITIALIZED) {
+ if (!file || file->type == HDFS_STREAM_UNINITIALIZED) {
errno = EBADF;
return -1;
}
- //The interface whose 'close' method to be called
- const char* interface = (file->type == INPUT) ?
+ interface = (file->type == HDFS_STREAM_INPUT) ?
HADOOP_ISTRM : HADOOP_OSTRM;
jthr = invokeMethod(env, NULL, INSTANCE, file->file, interface,
"close", "()V");
if (jthr) {
- const char *interfaceShortName = (file->type == INPUT) ?
+ interfaceShortName = (file->type == HDFS_STREAM_INPUT) ?
"FSDataInputStream" : "FSDataOutputStream";
ret = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
"%s#close", interfaceShortName);
@@ -1044,15 +1051,15 @@ int hdfsCloseFile(hdfsFS fs, hdfsFile fi
int hdfsExists(hdfsFS fs, const char *path)
{
JNIEnv *env = getJNIEnv();
- if (env == NULL) {
- errno = EINTERNAL;
- return -1;
- }
-
jobject jPath;
jvalue jVal;
jobject jFS = (jobject)fs;
jthrowable jthr;
+
+ if (env == NULL) {
+ errno = EINTERNAL;
+ return -1;
+ }
if (path == NULL) {
errno = EINVAL;
@@ -1088,13 +1095,13 @@ static int readPrepare(JNIEnv* env, hdfs
*jInputStream = (jobject)(f ? f->file : NULL);
//Sanity check
- if (!f || f->type == UNINITIALIZED) {
+ if (!f || f->type == HDFS_STREAM_UNINITIALIZED) {
errno = EBADF;
return -1;
}
//Error checking... make sure that this file is 'readable'
- if (f->type != INPUT) {
+ if (f->type != HDFS_STREAM_INPUT) {
fprintf(stderr, "Cannot read from a non-InputStream object!\n");
errno = EINVAL;
return -1;
@@ -1105,6 +1112,13 @@ static int readPrepare(JNIEnv* env, hdfs
tSize hdfsRead(hdfsFS fs, hdfsFile f, void* buffer, tSize length)
{
+ jobject jInputStream;
+ jbyteArray jbRarray;
+ jint noReadBytes = length;
+ jvalue jVal;
+ jthrowable jthr;
+ JNIEnv* env;
+
if (length == 0) {
return 0;
} else if (length < 0) {
@@ -1120,23 +1134,17 @@ tSize hdfsRead(hdfsFS fs, hdfsFile f, vo
// fis.read(bR);
//Get the JNIEnv* corresponding to current thread
- JNIEnv* env = getJNIEnv();
+ env = getJNIEnv();
if (env == NULL) {
errno = EINTERNAL;
return -1;
}
//Parameters
- jobject jInputStream;
if (readPrepare(env, fs, f, &jInputStream) == -1) {
return -1;
}
- jbyteArray jbRarray;
- jint noReadBytes = length;
- jvalue jVal;
- jthrowable jthr;
-
//Read the requisite bytes
jbRarray = (*env)->NewByteArray(env, length);
if (!jbRarray) {
@@ -1179,6 +1187,11 @@ tSize readDirect(hdfsFS fs, hdfsFile f,
// ByteBuffer bbuffer = ByteBuffer.allocateDirect(length) // wraps C buffer
// fis.read(bbuffer);
+ jobject jInputStream;
+ jvalue jVal;
+ jthrowable jthr;
+ jobject bb;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1186,16 +1199,12 @@ tSize readDirect(hdfsFS fs, hdfsFile f,
return -1;
}
- jobject jInputStream;
if (readPrepare(env, fs, f, &jInputStream) == -1) {
return -1;
}
- jvalue jVal;
- jthrowable jthr;
-
//Read the requisite bytes
- jobject bb = (*env)->NewDirectByteBuffer(env, buffer, length);
+ bb = (*env)->NewDirectByteBuffer(env, buffer, length);
if (bb == NULL) {
errno = printPendingExceptionAndFree(env, PRINT_EXC_ALL,
"readDirect: NewDirectByteBuffer");
@@ -1227,7 +1236,7 @@ tSize hdfsPread(hdfsFS fs, hdfsFile f, t
errno = EINVAL;
return -1;
}
- if (!f || f->type == UNINITIALIZED) {
+ if (!f || f->type == HDFS_STREAM_UNINITIALIZED) {
errno = EBADF;
return -1;
}
@@ -1239,7 +1248,7 @@ tSize hdfsPread(hdfsFS fs, hdfsFile f, t
}
//Error checking... make sure that this file is 'readable'
- if (f->type != INPUT) {
+ if (f->type != HDFS_STREAM_INPUT) {
fprintf(stderr, "Cannot read from a non-InputStream object!\n");
errno = EINVAL;
return -1;
@@ -1287,6 +1296,10 @@ tSize hdfsWrite(hdfsFS fs, hdfsFile f, c
// byte b[] = str.getBytes();
// fso.write(b);
+ jobject jOutputStream;
+ jbyteArray jbWarray;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1295,14 +1308,12 @@ tSize hdfsWrite(hdfsFS fs, hdfsFile f, c
}
//Sanity check
- if (!f || f->type == UNINITIALIZED) {
+ if (!f || f->type == HDFS_STREAM_UNINITIALIZED) {
errno = EBADF;
return -1;
}
- jobject jOutputStream = f->file;
- jbyteArray jbWarray;
- jthrowable jthr;
+ jOutputStream = f->file;
if (length < 0) {
errno = EINVAL;
@@ -1310,7 +1321,7 @@ tSize hdfsWrite(hdfsFS fs, hdfsFile f, c
}
//Error checking... make sure that this file is 'writable'
- if (f->type != OUTPUT) {
+ if (f->type != HDFS_STREAM_OUTPUT) {
fprintf(stderr, "Cannot write into a non-OutputStream object!\n");
errno = EINVAL;
return -1;
@@ -1355,6 +1366,9 @@ int hdfsSeek(hdfsFS fs, hdfsFile f, tOff
// JAVA EQUIVALENT
// fis.seek(pos);
+ jobject jInputStream;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1363,13 +1377,13 @@ int hdfsSeek(hdfsFS fs, hdfsFile f, tOff
}
//Sanity check
- if (!f || f->type != INPUT) {
+ if (!f || f->type != HDFS_STREAM_INPUT) {
errno = EBADF;
return -1;
}
- jobject jInputStream = f->file;
- jthrowable jthr = invokeMethod(env, NULL, INSTANCE, jInputStream,
+ jInputStream = f->file;
+ jthr = invokeMethod(env, NULL, INSTANCE, jInputStream,
HADOOP_ISTRM, "seek", "(J)V", desiredPos);
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1387,6 +1401,11 @@ tOffset hdfsTell(hdfsFS fs, hdfsFile f)
// JAVA EQUIVALENT
// pos = f.getPos();
+ jobject jStream;
+ const char *interface;
+ jvalue jVal;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1395,22 +1414,21 @@ tOffset hdfsTell(hdfsFS fs, hdfsFile f)
}
//Sanity check
- if (!f || f->type == UNINITIALIZED) {
+ if (!f || f->type == HDFS_STREAM_UNINITIALIZED) {
errno = EBADF;
return -1;
}
//Parameters
- jobject jStream = f->file;
- const char* interface = (f->type == INPUT) ?
+ jStream = f->file;
+ interface = (f->type == HDFS_STREAM_INPUT) ?
HADOOP_ISTRM : HADOOP_OSTRM;
- jvalue jVal;
- jthrowable jthr = invokeMethod(env, &jVal, INSTANCE, jStream,
+ jthr = invokeMethod(env, &jVal, INSTANCE, jStream,
interface, "getPos", "()J");
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
"hdfsTell: %s#getPos",
- ((f->type == INPUT) ? "FSDataInputStream" :
+ ((f->type == HDFS_STREAM_INPUT) ? "FSDataInputStream" :
"FSDataOutputStream"));
return -1;
}
@@ -1422,6 +1440,8 @@ int hdfsFlush(hdfsFS fs, hdfsFile f)
// JAVA EQUIVALENT
// fos.flush();
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1430,11 +1450,11 @@ int hdfsFlush(hdfsFS fs, hdfsFile f)
}
//Sanity check
- if (!f || f->type != OUTPUT) {
+ if (!f || f->type != HDFS_STREAM_OUTPUT) {
errno = EBADF;
return -1;
}
- jthrowable jthr = invokeMethod(env, NULL, INSTANCE, f->file,
+ jthr = invokeMethod(env, NULL, INSTANCE, f->file,
HADOOP_OSTRM, "flush", "()V");
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1446,6 +1466,9 @@ int hdfsFlush(hdfsFS fs, hdfsFile f)
int hdfsHFlush(hdfsFS fs, hdfsFile f)
{
+ jobject jOutputStream;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1454,13 +1477,13 @@ int hdfsHFlush(hdfsFS fs, hdfsFile f)
}
//Sanity check
- if (!f || f->type != OUTPUT) {
+ if (!f || f->type != HDFS_STREAM_OUTPUT) {
errno = EBADF;
return -1;
}
- jobject jOutputStream = f->file;
- jthrowable jthr = invokeMethod(env, NULL, INSTANCE, jOutputStream,
+ jOutputStream = f->file;
+ jthr = invokeMethod(env, NULL, INSTANCE, jOutputStream,
HADOOP_OSTRM, "hflush", "()V");
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1472,6 +1495,9 @@ int hdfsHFlush(hdfsFS fs, hdfsFile f)
int hdfsHSync(hdfsFS fs, hdfsFile f)
{
+ jobject jOutputStream;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1480,13 +1506,13 @@ int hdfsHSync(hdfsFS fs, hdfsFile f)
}
//Sanity check
- if (!f || f->type != OUTPUT) {
+ if (!f || f->type != HDFS_STREAM_OUTPUT) {
errno = EBADF;
return -1;
}
- jobject jOutputStream = f->file;
- jthrowable jthr = invokeMethod(env, NULL, INSTANCE, jOutputStream,
+ jOutputStream = f->file;
+ jthr = invokeMethod(env, NULL, INSTANCE, jOutputStream,
HADOOP_OSTRM, "hsync", "()V");
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1501,6 +1527,10 @@ int hdfsAvailable(hdfsFS fs, hdfsFile f)
// JAVA EQUIVALENT
// fis.available();
+ jobject jInputStream;
+ jvalue jVal;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1509,15 +1539,14 @@ int hdfsAvailable(hdfsFS fs, hdfsFile f)
}
//Sanity check
- if (!f || f->type != INPUT) {
+ if (!f || f->type != HDFS_STREAM_INPUT) {
errno = EBADF;
return -1;
}
//Parameters
- jobject jInputStream = f->file;
- jvalue jVal;
- jthrowable jthr = invokeMethod(env, &jVal, INSTANCE, jInputStream,
+ jInputStream = f->file;
+ jthr = invokeMethod(env, &jVal, INSTANCE, jInputStream,
HADOOP_ISTRM, "available", "()I");
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1527,20 +1556,13 @@ int hdfsAvailable(hdfsFS fs, hdfsFile f)
return jVal.i;
}
-static int hdfsCopyImpl(hdfsFS srcFS, const char* src, hdfsFS dstFS,
- const char* dst, jboolean deleteSource)
+static int hdfsCopyImpl(hdfsFS srcFS, const char *src, hdfsFS dstFS,
+ const char *dst, jboolean deleteSource)
{
//JAVA EQUIVALENT
// FileUtil#copy(srcFS, srcPath, dstFS, dstPath,
// deleteSource = false, conf)
- //Get the JNIEnv* corresponding to current thread
- JNIEnv* env = getJNIEnv();
- if (env == NULL) {
- errno = EINTERNAL;
- return -1;
- }
-
//Parameters
jobject jSrcFS = (jobject)srcFS;
jobject jDstFS = (jobject)dstFS;
@@ -1549,6 +1571,13 @@ static int hdfsCopyImpl(hdfsFS srcFS, co
jvalue jVal;
int ret;
+ //Get the JNIEnv* corresponding to current thread
+ JNIEnv* env = getJNIEnv();
+ if (env == NULL) {
+ errno = EINTERNAL;
+ return -1;
+ }
+
jthr = constructNewObjectOfPath(env, src, &jSrcPath);
if (jthr) {
ret = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1603,22 +1632,28 @@ done:
return 0;
}
-int hdfsCopy(hdfsFS srcFS, const char* src, hdfsFS dstFS, const char* dst)
+int hdfsCopy(hdfsFS srcFS, const char *src, hdfsFS dstFS, const char *dst)
{
return hdfsCopyImpl(srcFS, src, dstFS, dst, 0);
}
-int hdfsMove(hdfsFS srcFS, const char* src, hdfsFS dstFS, const char* dst)
+int hdfsMove(hdfsFS srcFS, const char *src, hdfsFS dstFS, const char *dst)
{
return hdfsCopyImpl(srcFS, src, dstFS, dst, 1);
}
-int hdfsDelete(hdfsFS fs, const char* path, int recursive)
+int hdfsDelete(hdfsFS fs, const char *path, int recursive)
{
// JAVA EQUIVALENT:
// Path p = new Path(path);
// bool retval = fs.delete(p, recursive);
+ jobject jFS = (jobject)fs;
+ jthrowable jthr;
+ jobject jPath;
+ jvalue jVal;
+ jboolean jRecursive;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1626,18 +1661,13 @@ int hdfsDelete(hdfsFS fs, const char* pa
return -1;
}
- jobject jFS = (jobject)fs;
- jthrowable jthr;
- jobject jPath;
- jvalue jVal;
-
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
"hdfsDelete(path=%s): constructNewObjectOfPath", path);
return -1;
}
- jboolean jRecursive = recursive ? JNI_TRUE : JNI_FALSE;
+ jRecursive = recursive ? JNI_TRUE : JNI_FALSE;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"delete", "(Lorg/apache/hadoop/fs/Path;Z)Z",
jPath, jRecursive);
@@ -1657,13 +1687,19 @@ int hdfsDelete(hdfsFS fs, const char* pa
-int hdfsRename(hdfsFS fs, const char* oldPath, const char* newPath)
+int hdfsRename(hdfsFS fs, const char *oldPath, const char *newPath)
{
// JAVA EQUIVALENT:
// Path old = new Path(oldPath);
// Path new = new Path(newPath);
// fs.rename(old, new);
+ jobject jFS = (jobject)fs;
+ jthrowable jthr;
+ jobject jOldPath = NULL, jNewPath = NULL;
+ int ret = -1;
+ jvalue jVal;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1671,12 +1707,6 @@ int hdfsRename(hdfsFS fs, const char* ol
return -1;
}
- jobject jFS = (jobject)fs;
- jthrowable jthr;
- jobject jOldPath = NULL, jNewPath = NULL;
- int ret = -1;
- jvalue jVal;
-
jthr = constructNewObjectOfPath(env, oldPath, &jOldPath );
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1721,13 +1751,6 @@ char* hdfsGetWorkingDirectory(hdfsFS fs,
// Path p = fs.getWorkingDirectory();
// return p.toString()
- //Get the JNIEnv* corresponding to current thread
- JNIEnv* env = getJNIEnv();
- if (env == NULL) {
- errno = EINTERNAL;
- return NULL;
- }
-
jobject jPath = NULL;
jstring jPathString = NULL;
jobject jFS = (jobject)fs;
@@ -1736,6 +1759,13 @@ char* hdfsGetWorkingDirectory(hdfsFS fs,
int ret;
const char *jPathChars = NULL;
+ //Get the JNIEnv* corresponding to current thread
+ JNIEnv* env = getJNIEnv();
+ if (env == NULL) {
+ errno = EINTERNAL;
+ return NULL;
+ }
+
//FileSystem#getWorkingDirectory()
jthr = invokeMethod(env, &jVal, INSTANCE, jFS,
HADOOP_FS, "getWorkingDirectory",
@@ -1794,11 +1824,15 @@ done:
-int hdfsSetWorkingDirectory(hdfsFS fs, const char* path)
+int hdfsSetWorkingDirectory(hdfsFS fs, const char *path)
{
// JAVA EQUIVALENT:
// fs.setWorkingDirectory(Path(path));
+ jobject jFS = (jobject)fs;
+ jthrowable jthr;
+ jobject jPath;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1806,10 +1840,6 @@ int hdfsSetWorkingDirectory(hdfsFS fs, c
return -1;
}
- jobject jFS = (jobject)fs;
- jthrowable jthr;
- jobject jPath;
-
//Create an object of org.apache.hadoop.fs.Path
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
@@ -1835,11 +1865,16 @@ int hdfsSetWorkingDirectory(hdfsFS fs, c
-int hdfsCreateDirectory(hdfsFS fs, const char* path)
+int hdfsCreateDirectory(hdfsFS fs, const char *path)
{
// JAVA EQUIVALENT:
// fs.mkdirs(new Path(path));
+ jobject jFS = (jobject)fs;
+ jobject jPath;
+ jthrowable jthr;
+ jvalue jVal;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1847,10 +1882,6 @@ int hdfsCreateDirectory(hdfsFS fs, const
return -1;
}
- jobject jFS = (jobject)fs;
- jobject jPath;
- jthrowable jthr;
-
//Create an object of org.apache.hadoop.fs.Path
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
@@ -1860,7 +1891,6 @@ int hdfsCreateDirectory(hdfsFS fs, const
}
//Create the directory
- jvalue jVal;
jVal.z = 0;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"mkdirs", "(Lorg/apache/hadoop/fs/Path;)Z",
@@ -1886,11 +1916,16 @@ int hdfsCreateDirectory(hdfsFS fs, const
}
-int hdfsSetReplication(hdfsFS fs, const char* path, int16_t replication)
+int hdfsSetReplication(hdfsFS fs, const char *path, int16_t replication)
{
// JAVA EQUIVALENT:
// fs.setReplication(new Path(path), replication);
+ jobject jFS = (jobject)fs;
+ jthrowable jthr;
+ jobject jPath;
+ jvalue jVal;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1898,11 +1933,7 @@ int hdfsSetReplication(hdfsFS fs, const
return -1;
}
- jobject jFS = (jobject)fs;
- jthrowable jthr;
-
//Create an object of org.apache.hadoop.fs.Path
- jobject jPath;
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -1911,7 +1942,6 @@ int hdfsSetReplication(hdfsFS fs, const
}
//Create the directory
- jvalue jVal;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"setReplication", "(Lorg/apache/hadoop/fs/Path;S)Z",
jPath, replication);
@@ -1932,11 +1962,17 @@ int hdfsSetReplication(hdfsFS fs, const
return 0;
}
-int hdfsChown(hdfsFS fs, const char* path, const char *owner, const char *group)
+int hdfsChown(hdfsFS fs, const char *path, const char *owner, const char *group)
{
// JAVA EQUIVALENT:
// fs.setOwner(path, owner, group)
+ jobject jFS = (jobject)fs;
+ jobject jPath = NULL;
+ jstring jOwner = NULL, jGroup = NULL;
+ jthrowable jthr;
+ int ret;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -1948,12 +1984,6 @@ int hdfsChown(hdfsFS fs, const char* pat
return 0;
}
- jobject jFS = (jobject)fs;
- jobject jPath = NULL;
- jstring jOwner = NULL, jGroup = NULL;
- jthrowable jthr;
- int ret;
-
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
ret = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -2001,12 +2031,17 @@ done:
return 0;
}
-int hdfsChmod(hdfsFS fs, const char* path, short mode)
+int hdfsChmod(hdfsFS fs, const char *path, short mode)
{
int ret;
// JAVA EQUIVALENT:
// fs.setPermission(path, FsPermission)
+ jthrowable jthr;
+ jobject jPath = NULL, jPermObj = NULL;
+ jobject jFS = (jobject)fs;
+ jshort jmode = mode;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -2014,12 +2049,7 @@ int hdfsChmod(hdfsFS fs, const char* pat
return -1;
}
- jthrowable jthr;
- jobject jPath = NULL, jPermObj = NULL;
- jobject jFS = (jobject)fs;
-
// construct jPerm = FsPermission.createImmutable(short mode);
- jshort jmode = mode;
jthr = constructNewObjectOfClass(env, &jPermObj,
HADOOP_FSPERM,"(S)V",jmode);
if (jthr) {
@@ -2061,11 +2091,16 @@ done:
return 0;
}
-int hdfsUtime(hdfsFS fs, const char* path, tTime mtime, tTime atime)
+int hdfsUtime(hdfsFS fs, const char *path, tTime mtime, tTime atime)
{
// JAVA EQUIVALENT:
// fs.setTimes(src, mtime, atime)
+
jthrowable jthr;
+ jobject jFS = (jobject)fs;
+ jobject jPath;
+ static const tTime NO_CHANGE = -1;
+ jlong jmtime, jatime;
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
@@ -2074,10 +2109,7 @@ int hdfsUtime(hdfsFS fs, const char* pat
return -1;
}
- jobject jFS = (jobject)fs;
-
//Create an object of org.apache.hadoop.fs.Path
- jobject jPath;
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
@@ -2085,9 +2117,8 @@ int hdfsUtime(hdfsFS fs, const char* pat
return -1;
}
- const tTime NO_CHANGE = -1;
- jlong jmtime = (mtime == NO_CHANGE) ? -1 : (mtime * (jlong)1000);
- jlong jatime = (atime == NO_CHANGE) ? -1 : (atime * (jlong)1000);
+ jmtime = (mtime == NO_CHANGE) ? -1 : (mtime * (jlong)1000);
+ jatime = (atime == NO_CHANGE) ? -1 : (atime * (jlong)1000);
jthr = invokeMethod(env, NULL, INSTANCE, jFS, HADOOP_FS,
"setTimes", JMETHOD3(JPARAM(HADOOP_PATH), "J", "J", JAVA_VOID),
@@ -2397,7 +2428,7 @@ struct hadoopRzBuffer* hadoopReadZero(hd
errno = EINTERNAL;
return NULL;
}
- if (file->type != INPUT) {
+ if (file->type != HDFS_STREAM_INPUT) {
fputs("Cannot read from a non-InputStream object!\n", stderr);
ret = EINVAL;
goto done;
@@ -2495,10 +2526,12 @@ void hadoopRzBufferFree(hdfsFile file, s
}
char***
-hdfsGetHosts(hdfsFS fs, const char* path, tOffset start, tOffset length)
+hdfsGetHosts(hdfsFS fs, const char *path, tOffset start, tOffset length)
{
// JAVA EQUIVALENT:
// fs.getFileBlockLoctions(new Path(path), start, length);
+
+ jobject jFS = (jobject)fs;
jthrowable jthr;
jobject jPath = NULL;
jobject jFileStatus = NULL;
@@ -2508,6 +2541,9 @@ hdfsGetHosts(hdfsFS fs, const char* path
char*** blockHosts = NULL;
int i, j, ret;
jsize jNumFileBlocks = 0;
+ jobject jFileBlock;
+ jsize jNumBlockHosts;
+ const char *hostName;
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
@@ -2516,8 +2552,6 @@ hdfsGetHosts(hdfsFS fs, const char* path
return NULL;
}
- jobject jFS = (jobject)fs;
-
//Create an object of org.apache.hadoop.fs.Path
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
@@ -2567,7 +2601,7 @@ hdfsGetHosts(hdfsFS fs, const char* path
//Now parse each block to get hostnames
for (i = 0; i < jNumFileBlocks; ++i) {
- jobject jFileBlock =
+ jFileBlock =
(*env)->GetObjectArrayElement(env, jBlockLocations, i);
if (!jFileBlock) {
ret = printPendingExceptionAndFree(env, PRINT_EXC_ALL,
@@ -2593,7 +2627,7 @@ hdfsGetHosts(hdfsFS fs, const char* path
goto done;
}
//Figure out no of hosts in jFileBlockHosts, and allocate the memory
- jsize jNumBlockHosts = (*env)->GetArrayLength(env, jFileBlockHosts);
+ jNumBlockHosts = (*env)->GetArrayLength(env, jFileBlockHosts);
blockHosts[i] = calloc(jNumBlockHosts + 1, sizeof(char*));
if (!blockHosts[i]) {
ret = ENOMEM;
@@ -2601,7 +2635,6 @@ hdfsGetHosts(hdfsFS fs, const char* path
}
//Now parse each hostname
- const char *hostName;
for (j = 0; j < jNumBlockHosts; ++j) {
jHost = (*env)->GetObjectArrayElement(env, jFileBlockHosts, j);
if (!jHost) {
@@ -2669,6 +2702,10 @@ tOffset hdfsGetDefaultBlockSize(hdfsFS f
// JAVA EQUIVALENT:
// fs.getDefaultBlockSize();
+ jobject jFS = (jobject)fs;
+ jvalue jVal;
+ jthrowable jthr;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -2676,11 +2713,7 @@ tOffset hdfsGetDefaultBlockSize(hdfsFS f
return -1;
}
- jobject jFS = (jobject)fs;
-
//FileSystem#getDefaultBlockSize()
- jvalue jVal;
- jthrowable jthr;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"getDefaultBlockSize", "()J");
if (jthr) {
@@ -2732,6 +2765,11 @@ tOffset hdfsGetCapacity(hdfsFS fs)
// FsStatus fss = fs.getStatus();
// return Fss.getCapacity();
+ jobject jFS = (jobject)fs;
+ jvalue jVal;
+ jthrowable jthr;
+ jobject fss;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -2739,11 +2777,7 @@ tOffset hdfsGetCapacity(hdfsFS fs)
return -1;
}
- jobject jFS = (jobject)fs;
-
//FileSystem#getStatus
- jvalue jVal;
- jthrowable jthr;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"getStatus", "()Lorg/apache/hadoop/fs/FsStatus;");
if (jthr) {
@@ -2751,7 +2785,7 @@ tOffset hdfsGetCapacity(hdfsFS fs)
"hdfsGetCapacity: FileSystem#getStatus");
return -1;
}
- jobject fss = (jobject)jVal.l;
+ fss = (jobject)jVal.l;
jthr = invokeMethod(env, &jVal, INSTANCE, fss, HADOOP_FSSTATUS,
"getCapacity", "()J");
destroyLocalReference(env, fss);
@@ -2771,6 +2805,11 @@ tOffset hdfsGetUsed(hdfsFS fs)
// FsStatus fss = fs.getStatus();
// return Fss.getUsed();
+ jobject jFS = (jobject)fs;
+ jvalue jVal;
+ jthrowable jthr;
+ jobject fss;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -2778,11 +2817,7 @@ tOffset hdfsGetUsed(hdfsFS fs)
return -1;
}
- jobject jFS = (jobject)fs;
-
//FileSystem#getStatus
- jvalue jVal;
- jthrowable jthr;
jthr = invokeMethod(env, &jVal, INSTANCE, jFS, HADOOP_FS,
"getStatus", "()Lorg/apache/hadoop/fs/FsStatus;");
if (jthr) {
@@ -2790,7 +2825,7 @@ tOffset hdfsGetUsed(hdfsFS fs)
"hdfsGetUsed: FileSystem#getStatus");
return -1;
}
- jobject fss = (jobject)jVal.l;
+ fss = (jobject)jVal.l;
jthr = invokeMethod(env, &jVal, INSTANCE, fss, HADOOP_FSSTATUS,
"getUsed", "()J");
destroyLocalReference(env, fss);
@@ -2814,6 +2849,9 @@ getFileInfoFromStat(JNIEnv *env, jobject
jstring jUserName = NULL;
jstring jGroupName = NULL;
jobject jPermission = NULL;
+ const char *cPathName;
+ const char *cUserName;
+ const char *cGroupName;
jthr = invokeMethod(env, &jVal, INSTANCE, jStat,
HADOOP_STAT, "isDir", "()Z");
@@ -2869,7 +2907,7 @@ getFileInfoFromStat(JNIEnv *env, jobject
if (jthr)
goto done;
jPathName = jVal.l;
- const char *cPathName =
+ cPathName =
(const char*) ((*env)->GetStringUTFChars(env, jPathName, NULL));
if (!cPathName) {
jthr = getPendingExceptionAndClear(env);
@@ -2882,7 +2920,7 @@ getFileInfoFromStat(JNIEnv *env, jobject
if (jthr)
goto done;
jUserName = jVal.l;
- const char* cUserName =
+ cUserName =
(const char*) ((*env)->GetStringUTFChars(env, jUserName, NULL));
if (!cUserName) {
jthr = getPendingExceptionAndClear(env);
@@ -2891,7 +2929,6 @@ getFileInfoFromStat(JNIEnv *env, jobject
fileInfo->mOwner = strdup(cUserName);
(*env)->ReleaseStringUTFChars(env, jUserName, cUserName);
- const char* cGroupName;
jthr = invokeMethod(env, &jVal, INSTANCE, jStat, HADOOP_STAT,
"getGroup", "()Ljava/lang/String;");
if (jthr)
@@ -2978,13 +3015,15 @@ getFileInfo(JNIEnv *env, jobject jFS, jo
-hdfsFileInfo* hdfsListDirectory(hdfsFS fs, const char* path, int *numEntries)
+hdfsFileInfo* hdfsListDirectory(hdfsFS fs, const char *path, int *numEntries)
{
// JAVA EQUIVALENT:
// Path p(path);
// Path []pathList = fs.listPaths(p)
// foreach path in pathList
// getFileInfo(path)
+
+ jobject jFS = (jobject)fs;
jthrowable jthr;
jobject jPath = NULL;
hdfsFileInfo *pathList = NULL;
@@ -2992,6 +3031,8 @@ hdfsFileInfo* hdfsListDirectory(hdfsFS f
jvalue jVal;
jsize jPathListSize = 0;
int ret;
+ jsize i;
+ jobject tmpStat;
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
@@ -3000,8 +3041,6 @@ hdfsFileInfo* hdfsListDirectory(hdfsFS f
return NULL;
}
- jobject jFS = (jobject)fs;
-
//Create an object of org.apache.hadoop.fs.Path
jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
@@ -3037,8 +3076,6 @@ hdfsFileInfo* hdfsListDirectory(hdfsFS f
}
//Save path information in pathList
- jsize i;
- jobject tmpStat;
for (i=0; i < jPathListSize; ++i) {
tmpStat = (*env)->GetObjectArrayElement(env, jPathList, i);
if (!tmpStat) {
@@ -3073,7 +3110,7 @@ done:
-hdfsFileInfo *hdfsGetPathInfo(hdfsFS fs, const char* path)
+hdfsFileInfo *hdfsGetPathInfo(hdfsFS fs, const char *path)
{
// JAVA EQUIVALENT:
// File f(path);
@@ -3082,6 +3119,11 @@ hdfsFileInfo *hdfsGetPathInfo(hdfsFS fs,
// fs.getLength(f)
// f.getPath()
+ jobject jFS = (jobject)fs;
+ jobject jPath;
+ jthrowable jthr;
+ hdfsFileInfo *fileInfo;
+
//Get the JNIEnv* corresponding to current thread
JNIEnv* env = getJNIEnv();
if (env == NULL) {
@@ -3089,17 +3131,13 @@ hdfsFileInfo *hdfsGetPathInfo(hdfsFS fs,
return NULL;
}
- jobject jFS = (jobject)fs;
-
//Create an object of org.apache.hadoop.fs.Path
- jobject jPath;
- jthrowable jthr = constructNewObjectOfPath(env, path, &jPath);
+ jthr = constructNewObjectOfPath(env, path, &jPath);
if (jthr) {
errno = printExceptionAndFree(env, jthr, PRINT_EXC_ALL,
"hdfsGetPathInfo(%s): constructNewObjectOfPath", path);
return NULL;
}
- hdfsFileInfo *fileInfo;
jthr = getFileInfo(env, jFS, jPath, &fileInfo);
destroyLocalReference(env, jPath);
if (jthr) {
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.c
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.c?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.c (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.c Wed Aug 20 01:34:29 2014
@@ -19,20 +19,18 @@
#include "config.h"
#include "exception.h"
#include "jni_helper.h"
+#include "platform.h"
+#include "common/htable.h"
+#include "os/mutexes.h"
+#include "os/thread_local_storage.h"
#include <stdio.h>
#include <string.h>
-static pthread_mutex_t hdfsHashMutex = PTHREAD_MUTEX_INITIALIZER;
-static pthread_mutex_t jvmMutex = PTHREAD_MUTEX_INITIALIZER;
-static volatile int hashTableInited = 0;
-
-#define LOCK_HASH_TABLE() pthread_mutex_lock(&hdfsHashMutex)
-#define UNLOCK_HASH_TABLE() pthread_mutex_unlock(&hdfsHashMutex)
-
+static struct htable *gClassRefHTable = NULL;
/** The Native return types that methods could return */
-#define VOID 'V'
+#define JVOID 'V'
#define JOBJECT 'L'
#define JARRAYOBJECT '['
#define JBOOLEAN 'Z'
@@ -51,40 +49,10 @@ static volatile int hashTableInited = 0;
*/
#define MAX_HASH_TABLE_ELEM 4096
-/** Key that allows us to retrieve thread-local storage */
-static pthread_key_t gTlsKey;
-
-/** nonzero if we succeeded in initializing gTlsKey. Protected by the jvmMutex */
-static int gTlsKeyInitialized = 0;
-
-/** Pthreads thread-local storage for each library thread. */
-struct hdfsTls {
- JNIEnv *env;
-};
-
/**
- * The function that is called whenever a thread with libhdfs thread local data
- * is destroyed.
- *
- * @param v The thread-local data
+ * Length of buffer for retrieving created JVMs. (We only ever create one.)
*/
-static void hdfsThreadDestructor(void *v)
-{
- struct hdfsTls *tls = v;
- JavaVM *vm;
- JNIEnv *env = tls->env;
- jint ret;
-
- ret = (*env)->GetJavaVM(env, &vm);
- if (ret) {
- fprintf(stderr, "hdfsThreadDestructor: GetJavaVM failed with "
- "error %d\n", ret);
- (*env)->ExceptionDescribe(env);
- } else {
- (*vm)->DetachCurrentThread(vm);
- }
- free(tls);
-}
+#define VM_BUF_LENGTH 1
void destroyLocalReference(JNIEnv *env, jobject jObject)
{
@@ -138,67 +106,6 @@ jthrowable newCStr(JNIEnv *env, jstring
return NULL;
}
-static int hashTableInit(void)
-{
- if (!hashTableInited) {
- LOCK_HASH_TABLE();
- if (!hashTableInited) {
- if (hcreate(MAX_HASH_TABLE_ELEM) == 0) {
- fprintf(stderr, "error creating hashtable, <%d>: %s\n",
- errno, strerror(errno));
- UNLOCK_HASH_TABLE();
- return 0;
- }
- hashTableInited = 1;
- }
- UNLOCK_HASH_TABLE();
- }
- return 1;
-}
-
-
-static int insertEntryIntoTable(const char *key, void *data)
-{
- ENTRY e, *ep;
- if (key == NULL || data == NULL) {
- return 0;
- }
- if (! hashTableInit()) {
- return -1;
- }
- e.data = data;
- e.key = (char*)key;
- LOCK_HASH_TABLE();
- ep = hsearch(e, ENTER);
- UNLOCK_HASH_TABLE();
- if (ep == NULL) {
- fprintf(stderr, "warn adding key (%s) to hash table, <%d>: %s\n",
- key, errno, strerror(errno));
- }
- return 0;
-}
-
-
-
-static void* searchEntryFromTable(const char *key)
-{
- ENTRY e,*ep;
- if (key == NULL) {
- return NULL;
- }
- hashTableInit();
- e.key = (char*)key;
- LOCK_HASH_TABLE();
- ep = hsearch(e, FIND);
- UNLOCK_HASH_TABLE();
- if (ep != NULL) {
- return ep->data;
- }
- return NULL;
-}
-
-
-
jthrowable invokeMethod(JNIEnv *env, jvalue *retval, MethType methType,
jobject instObj, const char *className,
const char *methName, const char *methSignature, ...)
@@ -235,7 +142,7 @@ jthrowable invokeMethod(JNIEnv *env, jva
}
retval->l = jobj;
}
- else if (returnType == VOID) {
+ else if (returnType == JVOID) {
if (methType == STATIC) {
(*env)->CallStaticVoidMethodV(env, cls, mid, args);
}
@@ -325,11 +232,11 @@ jthrowable methodIdFromClass(const char
{
jclass cls;
jthrowable jthr;
+ jmethodID mid = 0;
jthr = globalClassReference(className, env, &cls);
if (jthr)
return jthr;
- jmethodID mid = 0;
jthr = validateMethodType(env, methType);
if (jthr)
return jthr;
@@ -350,25 +257,50 @@ jthrowable methodIdFromClass(const char
jthrowable globalClassReference(const char *className, JNIEnv *env, jclass *out)
{
- jclass clsLocalRef;
- jclass cls = searchEntryFromTable(className);
- if (cls) {
- *out = cls;
- return NULL;
+ jthrowable jthr = NULL;
+ jclass local_clazz = NULL;
+ jclass clazz = NULL;
+ int ret;
+
+ mutexLock(&hdfsHashMutex);
+ if (!gClassRefHTable) {
+ gClassRefHTable = htable_alloc(MAX_HASH_TABLE_ELEM, ht_hash_string,
+ ht_compare_string);
+ if (!gClassRefHTable) {
+ jthr = newRuntimeError(env, "htable_alloc failed\n");
+ goto done;
+ }
}
- clsLocalRef = (*env)->FindClass(env,className);
- if (clsLocalRef == NULL) {
- return getPendingExceptionAndClear(env);
+ clazz = htable_get(gClassRefHTable, className);
+ if (clazz) {
+ *out = clazz;
+ goto done;
}
- cls = (*env)->NewGlobalRef(env, clsLocalRef);
- if (cls == NULL) {
- (*env)->DeleteLocalRef(env, clsLocalRef);
- return getPendingExceptionAndClear(env);
+ local_clazz = (*env)->FindClass(env,className);
+ if (!local_clazz) {
+ jthr = getPendingExceptionAndClear(env);
+ goto done;
}
- (*env)->DeleteLocalRef(env, clsLocalRef);
- insertEntryIntoTable(className, cls);
- *out = cls;
- return NULL;
+ clazz = (*env)->NewGlobalRef(env, local_clazz);
+ if (!clazz) {
+ jthr = getPendingExceptionAndClear(env);
+ goto done;
+ }
+ ret = htable_put(gClassRefHTable, (void*)className, clazz);
+ if (ret) {
+ jthr = newRuntimeError(env, "htable_put failed with error "
+ "code %d\n", ret);
+ goto done;
+ }
+ *out = clazz;
+ jthr = NULL;
+done:
+ mutexUnlock(&hdfsHashMutex);
+ (*env)->DeleteLocalRef(env, local_clazz);
+ if (jthr && clazz) {
+ (*env)->DeleteGlobalRef(env, clazz);
+ }
+ return jthr;
}
jthrowable classNameOfObject(jobject jobj, JNIEnv *env, char **name)
@@ -436,14 +368,24 @@ done:
*/
static JNIEnv* getGlobalJNIEnv(void)
{
- const jsize vmBufLength = 1;
- JavaVM* vmBuf[vmBufLength];
+ JavaVM* vmBuf[VM_BUF_LENGTH];
JNIEnv *env;
jint rv = 0;
jint noVMs = 0;
jthrowable jthr;
+ char *hadoopClassPath;
+ const char *hadoopClassPathVMArg = "-Djava.class.path=";
+ size_t optHadoopClassPathLen;
+ char *optHadoopClassPath;
+ int noArgs = 1;
+ char *hadoopJvmArgs;
+ char jvmArgDelims[] = " ";
+ char *str, *token, *savePtr;
+ JavaVMInitArgs vm_args;
+ JavaVM *vm;
+ JavaVMOption *options;
- rv = JNI_GetCreatedJavaVMs(&(vmBuf[0]), vmBufLength, &noVMs);
+ rv = JNI_GetCreatedJavaVMs(&(vmBuf[0]), VM_BUF_LENGTH, &noVMs);
if (rv != 0) {
fprintf(stderr, "JNI_GetCreatedJavaVMs failed with error: %d\n", rv);
return NULL;
@@ -451,23 +393,19 @@ static JNIEnv* getGlobalJNIEnv(void)
if (noVMs == 0) {
//Get the environment variables for initializing the JVM
- char *hadoopClassPath = getenv("CLASSPATH");
+ hadoopClassPath = getenv("CLASSPATH");
if (hadoopClassPath == NULL) {
fprintf(stderr, "Environment variable CLASSPATH not set!\n");
return NULL;
}
- char *hadoopClassPathVMArg = "-Djava.class.path=";
- size_t optHadoopClassPathLen = strlen(hadoopClassPath) +
+ optHadoopClassPathLen = strlen(hadoopClassPath) +
strlen(hadoopClassPathVMArg) + 1;
- char *optHadoopClassPath = malloc(sizeof(char)*optHadoopClassPathLen);
+ optHadoopClassPath = malloc(sizeof(char)*optHadoopClassPathLen);
snprintf(optHadoopClassPath, optHadoopClassPathLen,
"%s%s", hadoopClassPathVMArg, hadoopClassPath);
// Determine the # of LIBHDFS_OPTS args
- int noArgs = 1;
- char *hadoopJvmArgs = getenv("LIBHDFS_OPTS");
- char jvmArgDelims[] = " ";
- char *str, *token, *savePtr;
+ hadoopJvmArgs = getenv("LIBHDFS_OPTS");
if (hadoopJvmArgs != NULL) {
hadoopJvmArgs = strdup(hadoopJvmArgs);
for (noArgs = 1, str = hadoopJvmArgs; ; noArgs++, str = NULL) {
@@ -480,7 +418,12 @@ static JNIEnv* getGlobalJNIEnv(void)
}
// Now that we know the # args, populate the options array
- JavaVMOption options[noArgs];
+ options = calloc(noArgs, sizeof(JavaVMOption));
+ if (!options) {
+ fputs("Call to calloc failed\n", stderr);
+ free(optHadoopClassPath);
+ return NULL;
+ }
options[0].optionString = optHadoopClassPath;
hadoopJvmArgs = getenv("LIBHDFS_OPTS");
if (hadoopJvmArgs != NULL) {
@@ -495,8 +438,6 @@ static JNIEnv* getGlobalJNIEnv(void)
}
//Create the VM
- JavaVMInitArgs vm_args;
- JavaVM *vm;
vm_args.version = JNI_VERSION_1_2;
vm_args.options = options;
vm_args.nOptions = noArgs;
@@ -508,6 +449,7 @@ static JNIEnv* getGlobalJNIEnv(void)
free(hadoopJvmArgs);
}
free(optHadoopClassPath);
+ free(options);
if (rv != 0) {
fprintf(stderr, "Call to JNI_CreateJavaVM failed "
@@ -523,7 +465,7 @@ static JNIEnv* getGlobalJNIEnv(void)
}
else {
//Attach this thread to the VM
- JavaVM* vm = vmBuf[0];
+ vm = vmBuf[0];
rv = (*vm)->AttachCurrentThread(vm, (void*)&env, 0);
if (rv != 0) {
fprintf(stderr, "Call to AttachCurrentThread "
@@ -557,54 +499,27 @@ static JNIEnv* getGlobalJNIEnv(void)
JNIEnv* getJNIEnv(void)
{
JNIEnv *env;
- struct hdfsTls *tls;
- int ret;
-
-#ifdef HAVE_BETTER_TLS
- static __thread struct hdfsTls *quickTls = NULL;
- if (quickTls)
- return quickTls->env;
-#endif
- pthread_mutex_lock(&jvmMutex);
- if (!gTlsKeyInitialized) {
- ret = pthread_key_create(&gTlsKey, hdfsThreadDestructor);
- if (ret) {
- pthread_mutex_unlock(&jvmMutex);
- fprintf(stderr, "getJNIEnv: pthread_key_create failed with "
- "error %d\n", ret);
- return NULL;
- }
- gTlsKeyInitialized = 1;
- }
- tls = pthread_getspecific(gTlsKey);
- if (tls) {
- pthread_mutex_unlock(&jvmMutex);
- return tls->env;
+ THREAD_LOCAL_STORAGE_GET_QUICK();
+ mutexLock(&jvmMutex);
+ if (threadLocalStorageGet(&env)) {
+ mutexUnlock(&jvmMutex);
+ return NULL;
+ }
+ if (env) {
+ mutexUnlock(&jvmMutex);
+ return env;
}
env = getGlobalJNIEnv();
- pthread_mutex_unlock(&jvmMutex);
+ mutexUnlock(&jvmMutex);
if (!env) {
- fprintf(stderr, "getJNIEnv: getGlobalJNIEnv failed\n");
- return NULL;
+ fprintf(stderr, "getJNIEnv: getGlobalJNIEnv failed\n");
+ return NULL;
}
- tls = calloc(1, sizeof(struct hdfsTls));
- if (!tls) {
- fprintf(stderr, "getJNIEnv: OOM allocating %zd bytes\n",
- sizeof(struct hdfsTls));
- return NULL;
- }
- tls->env = env;
- ret = pthread_setspecific(gTlsKey, tls);
- if (ret) {
- fprintf(stderr, "getJNIEnv: pthread_setspecific failed with "
- "error code %d\n", ret);
- hdfsThreadDestructor(tls);
- return NULL;
+ if (threadLocalStorageSet(env)) {
+ return NULL;
}
-#ifdef HAVE_BETTER_TLS
- quickTls = tls;
-#endif
+ THREAD_LOCAL_STORAGE_SET_QUICK(env);
return env;
}
Modified: hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.h
URL: http://svn.apache.org/viewvc/hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.h?rev=1619019&r1=1619018&r2=1619019&view=diff
==============================================================================
--- hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.h (original)
+++ hadoop/common/branches/YARN-1051/hadoop-hdfs-project/hadoop-hdfs/src/main/native/libhdfs/jni_helper.h Wed Aug 20 01:34:29 2014
@@ -24,8 +24,6 @@
#include <stdlib.h>
#include <stdarg.h>
-#include <search.h>
-#include <pthread.h>
#include <errno.h>
#define PATH_SEPARATOR ':'