codekingpro/portable-devtools
114k
1from __future__ import annotations2 3import concurrent.futures4from pathlib import Path5from typing import Iterator, Literal, Optional, Sequence, Union6 7from langchain_core.documents import Document8 9from langchain_community.document_loaders.base import BaseBlobParser10from langchain_community.document_loaders.blob_loaders import (11 BlobLoader,12 FileSystemBlobLoader,13)14from langchain_community.document_loaders.generic import GenericLoader15from langchain_community.document_loaders.parsers.registry import get_parser16 17_PathLike = Union[str, Path]18 19DEFAULT = Literal["default"]20 21 22class ConcurrentLoader(GenericLoader):23 """Load and pars Documents concurrently."""24 25 def __init__(26 self,27 blob_loader: BlobLoader,28 blob_parser: BaseBlobParser,29 num_workers: int = 4,30 ) -> None:31 super().__init__(blob_loader, blob_parser)32 self.num_workers = num_workers33 34 def lazy_load(35 self,36 ) -> Iterator[Document]:37 """Load documents lazily with concurrent parsing."""38 with concurrent.futures.ThreadPoolExecutor(39 max_workers=self.num_workers40 ) as executor:41 futures = {42 executor.submit(self.blob_parser.lazy_parse, blob)43 for blob in self.blob_loader.yield_blobs()44 }45 for future in concurrent.futures.as_completed(futures):46 yield from future.result()47 48 @classmethod49 def from_filesystem(50 cls,51 path: _PathLike,52 *,53 glob: str = "**/[!.]*",54 exclude: Sequence[str] = (),55 suffixes: Optional[Sequence[str]] = None,56 show_progress: bool = False,57 parser: Union[DEFAULT, BaseBlobParser] = "default",58 num_workers: int = 4,59 parser_kwargs: Optional[dict] = None,60 ) -> ConcurrentLoader:61 """Create a concurrent generic document loader using a filesystem blob loader.62 63 Args:64 path: The path to the directory to load documents from.65 glob: The glob pattern to use to find documents.66 suffixes: The suffixes to use to filter documents. If None, all files67 matching the glob will be loaded.68 exclude: A list of patterns to exclude from the loader.69 show_progress: Whether to show a progress bar or not (requires tqdm).70 Proxies to the file system loader.71 parser: A blob parser which knows how to parse blobs into documents72 num_workers: Max number of concurrent workers to use.73 parser_kwargs: Keyword arguments to pass to the parser.74 """75 blob_loader = FileSystemBlobLoader(76 path,77 glob=glob,78 exclude=exclude,79 suffixes=suffixes,80 show_progress=show_progress,81 )82 if isinstance(parser, str):83 if parser == "default" and cls.get_parser != GenericLoader.get_parser:84 # There is an implementation of get_parser on the class, use it.85 blob_parser = cls.get_parser(**(parser_kwargs or {}))86 else:87 blob_parser = get_parser(parser)88 else:89 blob_parser = parser90 return cls(blob_loader, blob_parser, num_workers=num_workers)91 