codekingpro/portable-devtools
114k
1# https://hadoop.apache.org/docs/r1.0.4/webhdfs.html2 3import logging4import os5import secrets6import shutil7import tempfile8import uuid9from contextlib import suppress10from datetime import datetime11from urllib.parse import quote12 13import requests14 15from ..spec import AbstractBufferedFile, AbstractFileSystem16from ..utils import infer_storage_options, tokenize17 18logger = logging.getLogger("webhdfs")19 20 21class WebHDFS(AbstractFileSystem):22 """23 Interface to HDFS over HTTP using the WebHDFS API. Supports also HttpFS gateways.24 25 Four auth mechanisms are supported:26 27 insecure: no auth is done, and the user is assumed to be whoever they28 say they are (parameter ``user``), or a predefined value such as29 "dr.who" if not given30 spnego: when kerberos authentication is enabled, auth is negotiated by31 requests_kerberos https://github.com/requests/requests-kerberos .32 This establishes a session based on existing kinit login and/or33 specified principal/password; parameters are passed with ``kerb_kwargs``34 token: uses an existing Hadoop delegation token from another secured35 service. Indeed, this client can also generate such tokens when36 not insecure. Note that tokens expire, but can be renewed (by a37 previously specified user) and may allow for proxying.38 basic-auth: used when both parameter ``user`` and parameter ``password``39 are provided.40 41 """42 43 tempdir = str(tempfile.gettempdir())44 protocol = "webhdfs", "webHDFS"45 46 def __init__(47 self,48 host,49 port=50070,50 kerberos=False,51 token=None,52 user=None,53 password=None,54 proxy_to=None,55 kerb_kwargs=None,56 data_proxy=None,57 use_https=False,58 session_cert=None,59 session_verify=True,60 **kwargs,61 ):62 """63 Parameters64 ----------65 host: str66 Name-node address67 port: int68 Port for webHDFS69 kerberos: bool70 Whether to authenticate with kerberos for this connection71 token: str or None72 If given, use this token on every call to authenticate. A user73 and user-proxy may be encoded in the token and should not be also74 given75 user: str or None76 If given, assert the user name to connect with77 password: str or None78 If given, assert the password to use for basic auth. If password79 is provided, user must be provided also80 proxy_to: str or None81 If given, the user has the authority to proxy, and this value is82 the user in who's name actions are taken83 kerb_kwargs: dict84 Any extra arguments for HTTPKerberosAuth, see85 `<https://github.com/requests/requests-kerberos/blob/master/requests_kerberos/kerberos_.py>`_86 data_proxy: dict, callable or None87 If given, map data-node addresses. This can be necessary if the88 HDFS cluster is behind a proxy, running on Docker or otherwise has89 a mismatch between the host-names given by the name-node and the90 address by which to refer to them from the client. If a dict,91 maps host names ``host->data_proxy[host]``; if a callable, full92 URLs are passed, and function must conform to93 ``url->data_proxy(url)``.94 use_https: bool95 Whether to connect to the Name-node using HTTPS instead of HTTP96 session_cert: str or Tuple[str, str] or None97 Path to a certificate file, or tuple of (cert, key) files to use98 for the requests.Session99 session_verify: str, bool or None100 Path to a certificate file to use for verifying the requests.Session.101 kwargs102 """103 if self._cached:104 return105 super().__init__(**kwargs)106 self.url = f"{'https' if use_https else 'http'}://{host}:{port}/webhdfs/v1"107 self.kerb = kerberos108 self.kerb_kwargs = kerb_kwargs or {}109 self.pars = {}110 self.proxy = data_proxy or {}111 if token is not None:112 if user is not None or proxy_to is not None:113 raise ValueError(114 "If passing a delegation token, must not set "115 "user or proxy_to, as these are encoded in the"116 " token"117 )118 self.pars["delegation"] = token119 self.user = user120 self.password = password121 122 if password is not None:123 if user is None:124 raise ValueError(125 "If passing a password, the user must also be"126 "set in order to set up the basic-auth"127 )128 else:129 if user is not None:130 self.pars["user.name"] = user131 132 if proxy_to is not None:133 self.pars["doas"] = proxy_to134 if kerberos and user is not None:135 raise ValueError(136 "If using Kerberos auth, do not specify the "137 "user, this is handled by kinit."138 )139 140 self.session_cert = session_cert141 self.session_verify = session_verify142 143 self._connect()144 145 self._fsid = f"webhdfs_{tokenize(host, port)}"146 147 @property148 def fsid(self):149 return self._fsid150 151 def _connect(self):152 self.session = requests.Session()153 154 if self.session_cert:155 self.session.cert = self.session_cert156 157 self.session.verify = self.session_verify158 159 if self.kerb:160 from requests_kerberos import HTTPKerberosAuth161 162 self.session.auth = HTTPKerberosAuth(**self.kerb_kwargs)163 164 if self.user is not None and self.password is not None:165 from requests.auth import HTTPBasicAuth166 167 self.session.auth = HTTPBasicAuth(self.user, self.password)168 169 def _call(self, op, method="get", path=None, data=None, redirect=True, **kwargs):170 path = self._strip_protocol(path) if path is not None else ""171 url = self._apply_proxy(self.url + quote(path, safe="/="))172 args = kwargs.copy()173 args.update(self.pars)174 args["op"] = op.upper()175 logger.debug("sending %s with %s", url, method)176 out = self.session.request(177 method=method.upper(),178 url=url,179 params=args,180 data=data,181 allow_redirects=redirect,182 )183 if out.status_code in [400, 401, 403, 404, 500]:184 try:185 err = out.json()186 msg = err["RemoteException"]["message"]187 exp = err["RemoteException"]["exception"]188 except (ValueError, KeyError):189 pass190 else:191 if exp in ["IllegalArgumentException", "UnsupportedOperationException"]:192 raise ValueError(msg)193 elif exp in ["SecurityException", "AccessControlException"]:194 raise PermissionError(msg)195 elif exp in ["FileNotFoundException"]:196 raise FileNotFoundError(msg)197 else:198 raise RuntimeError(msg)199 out.raise_for_status()200 return out201 202 def _open(203 self,204 path,205 mode="rb",206 block_size=None,207 autocommit=True,208 replication=None,209 permissions=None,210 **kwargs,211 ):212 """213 214 Parameters215 ----------216 path: str217 File location218 mode: str219 'rb', 'wb', etc.220 block_size: int221 Client buffer size for read-ahead or write buffer222 autocommit: bool223 If False, writes to temporary file that only gets put in final224 location upon commit225 replication: int226 Number of copies of file on the cluster, write mode only227 permissions: str or int228 posix permissions, write mode only229 kwargs230 231 Returns232 -------233 WebHDFile instance234 """235 block_size = block_size or self.blocksize236 return WebHDFile(237 self,238 path,239 mode=mode,240 block_size=block_size,241 tempdir=self.tempdir,242 autocommit=autocommit,243 replication=replication,244 permissions=permissions,245 )246 247 @staticmethod248 def _process_info(info):249 info["type"] = info["type"].lower()250 info["size"] = info["length"]251 return info252 253 @classmethod254 def _strip_protocol(cls, path):255 return infer_storage_options(path)["path"]256 257 @staticmethod258 def _get_kwargs_from_urls(urlpath):259 out = infer_storage_options(urlpath)260 out.pop("path", None)261 out.pop("protocol", None)262 if "username" in out:263 out["user"] = out.pop("username")264 return out265 266 def info(self, path):267 out = self._call("GETFILESTATUS", path=path)268 info = out.json()["FileStatus"]269 info["name"] = path270 return self._process_info(info)271 272 def created(self, path):273 """Return the created timestamp of a file as a datetime.datetime"""274 # The API does not provide creation time, so we use modification time275 info = self.info(path)276 mtime = info.get("modificationTime", None)277 if mtime is not None:278 return datetime.fromtimestamp(mtime / 1000)279 raise RuntimeError("Could not retrieve creation time (modification time).")280 281 def modified(self, path):282 """Return the modified timestamp of a file as a datetime.datetime"""283 info = self.info(path)284 mtime = info.get("modificationTime", None)285 if mtime is not None:286 return datetime.fromtimestamp(mtime / 1000)287 raise RuntimeError("Could not retrieve modification time.")288 289 def ls(self, path, detail=False, **kwargs):290 out = self._call("LISTSTATUS", path=path)291 infos = out.json()["FileStatuses"]["FileStatus"]292 for info in infos:293 self._process_info(info)294 info["name"] = path.rstrip("/") + "/" + info["pathSuffix"]295 if detail:296 return sorted(infos, key=lambda i: i["name"])297 else:298 return sorted(info["name"] for info in infos)299 300 def content_summary(self, path):301 """Total numbers of files, directories and bytes under path"""302 out = self._call("GETCONTENTSUMMARY", path=path)303 return out.json()["ContentSummary"]304 305 def ukey(self, path):306 """Checksum info of file, giving method and result"""307 out = self._call("GETFILECHECKSUM", path=path, redirect=False)308 if "Location" in out.headers:309 location = self._apply_proxy(out.headers["Location"])310 out2 = self.session.get(location)311 out2.raise_for_status()312 return out2.json()["FileChecksum"]313 else:314 out.raise_for_status()315 return out.json()["FileChecksum"]316 317 def home_directory(self):318 """Get user's home directory"""319 out = self._call("GETHOMEDIRECTORY")320 return out.json()["Path"]321 322 def get_delegation_token(self, renewer=None):323 """Retrieve token which can give the same authority to other uses324 325 Parameters326 ----------327 renewer: str or None328 User who may use this token; if None, will be current user329 """330 if renewer:331 out = self._call("GETDELEGATIONTOKEN", renewer=renewer)332 else:333 out = self._call("GETDELEGATIONTOKEN")334 t = out.json()["Token"]335 if t is None:336 raise ValueError("No token available for this user/security context")337 return t["urlString"]338 339 def renew_delegation_token(self, token):340 """Make token live longer. Returns new expiry time"""341 out = self._call("RENEWDELEGATIONTOKEN", method="put", token=token)342 return out.json()["long"]343 344 def cancel_delegation_token(self, token):345 """Stop the token from being useful"""346 self._call("CANCELDELEGATIONTOKEN", method="put", token=token)347 348 def chmod(self, path, mod):349 """Set the permission at path350 351 Parameters352 ----------353 path: str354 location to set (file or directory)355 mod: str or int356 posix epresentation or permission, give as oct string, e.g, '777'357 or 0o777358 """359 self._call("SETPERMISSION", method="put", path=path, permission=mod)360 361 def chown(self, path, owner=None, group=None):362 """Change owning user and/or group"""363 kwargs = {}364 if owner is not None:365 kwargs["owner"] = owner366 if group is not None:367 kwargs["group"] = group368 self._call("SETOWNER", method="put", path=path, **kwargs)369 370 def set_replication(self, path, replication):371 """372 Set file replication factor373 374 Parameters375 ----------376 path: str377 File location (not for directories)378 replication: int379 Number of copies of file on the cluster. Should be smaller than380 number of data nodes; normally 3 on most systems.381 """382 self._call("SETREPLICATION", path=path, method="put", replication=replication)383 384 def mkdir(self, path, **kwargs):385 self._call("MKDIRS", method="put", path=path)386 387 def makedirs(self, path, exist_ok=False):388 if exist_ok is False and self.exists(path):389 raise FileExistsError(path)390 self.mkdir(path)391 392 def mv(self, path1, path2, **kwargs):393 self._call("RENAME", method="put", path=path1, destination=path2)394 395 def rm(self, path, recursive=False, **kwargs):396 self._call(397 "DELETE",398 method="delete",399 path=path,400 recursive="true" if recursive else "false",401 )402 403 def rm_file(self, path, **kwargs):404 self.rm(path)405 406 def cp_file(self, lpath, rpath, **kwargs):407 with self.open(lpath) as lstream:408 tmp_fname = "/".join([self._parent(rpath), f".tmp.{secrets.token_hex(16)}"])409 # Perform an atomic copy (stream to a temporary file and410 # move it to the actual destination).411 try:412 with self.open(tmp_fname, "wb") as rstream:413 shutil.copyfileobj(lstream, rstream)414 self.mv(tmp_fname, rpath)415 except BaseException:416 with suppress(FileNotFoundError):417 self.rm(tmp_fname)418 raise419 420 def _apply_proxy(self, location):421 if self.proxy and callable(self.proxy):422 location = self.proxy(location)423 elif self.proxy:424 # as a dict425 for k, v in self.proxy.items():426 location = location.replace(k, v, 1)427 return location428 429 430class WebHDFile(AbstractBufferedFile):431 """A file living in HDFS over webHDFS"""432 433 def __init__(self, fs, path, **kwargs):434 super().__init__(fs, path, **kwargs)435 kwargs = kwargs.copy()436 if kwargs.get("permissions", None) is None:437 kwargs.pop("permissions", None)438 if kwargs.get("replication", None) is None:439 kwargs.pop("replication", None)440 self.permissions = kwargs.pop("permissions", 511)441 tempdir = kwargs.pop("tempdir")442 if kwargs.pop("autocommit", False) is False:443 self.target = self.path444 self.path = os.path.join(tempdir, str(uuid.uuid4()))445 446 def _upload_chunk(self, final=False):447 """Write one part of a multi-block file upload448 449 Parameters450 ==========451 final: bool452 This is the last block, so should complete file, if453 self.autocommit is True.454 """455 out = self.fs.session.post(456 self.location,457 data=self.buffer.getvalue(),458 headers={"content-type": "application/octet-stream"},459 )460 out.raise_for_status()461 return True462 463 def _initiate_upload(self):464 """Create remote file/upload"""465 kwargs = self.kwargs.copy()466 if "a" in self.mode:467 op, method = "APPEND", "POST"468 else:469 op, method = "CREATE", "PUT"470 kwargs["overwrite"] = "true"471 out = self.fs._call(op, method, self.path, redirect=False, **kwargs)472 location = self.fs._apply_proxy(out.headers["Location"])473 if "w" in self.mode:474 # create empty file to append to475 out2 = self.fs.session.put(476 location, headers={"content-type": "application/octet-stream"}477 )478 out2.raise_for_status()479 # after creating empty file, change location to append to480 out2 = self.fs._call("APPEND", "POST", self.path, redirect=False, **kwargs)481 self.location = self.fs._apply_proxy(out2.headers["Location"])482 483 def _fetch_range(self, start, end):484 start = max(start, 0)485 end = min(self.size, end)486 if start >= end or start >= self.size:487 return b""488 out = self.fs._call(489 "OPEN", path=self.path, offset=start, length=end - start, redirect=False490 )491 out.raise_for_status()492 if "Location" in out.headers:493 location = out.headers["Location"]494 out2 = self.fs.session.get(self.fs._apply_proxy(location))495 return out2.content496 else:497 return out.content498 499 def commit(self):500 self.fs.mv(self.path, self.target)501 502 def discard(self):503 self.fs.rm(self.path)504 