Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
webhdfs.py504 linesDownload Raw Back to implementations
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 
codekingpro/portable-devtools · Team Ai