Storage for file system.
| 68 | raise ValueError(f"Unsupported file type: {self.cache_type}, output file should end with json, jsonl, csv, parquet, pickle") |
| 69 | |
| 70 | class FileStorage(DataFlowStorage): |
| 71 | """ |
| 72 | Storage for file system. |
| 73 | """ |
| 74 | def __init__( |
| 75 | self, |
| 76 | first_entry_file_name: str, |
| 77 | cache_path:str="./cache", |
| 78 | file_name_prefix:str="dataflow_cache", |
| 79 | cache_type:Literal["json", "jsonl", "csv", "parquet", "pickle"] = "jsonl" |
| 80 | ): |
| 81 | self.first_entry_file_name = first_entry_file_name |
| 82 | self.cache_path = cache_path |
| 83 | self.file_name_prefix = file_name_prefix |
| 84 | self.cache_type = cache_type |
| 85 | self.operator_step = -1 |
| 86 | self.logger = get_logger() |
| 87 | |
| 88 | def _get_cache_file_path(self, step) -> str: |
| 89 | if step == -1: |
| 90 | self.logger.error("You must call storage.step() before reading or writing data. Please call storage.step() first for each operator step.") |
| 91 | raise ValueError("You must call storage.step() before reading or writing data. Please call storage.step() first for each operator step.") |
| 92 | if step == 0: |
| 93 | # If it's the first step, use the first entry file name |
| 94 | return os.path.join(self.first_entry_file_name) |
| 95 | else: |
| 96 | return os.path.join(self.cache_path, f"{self.file_name_prefix}_step{step}.{self.cache_type}") |
| 97 | |
| 98 | def step(self): |
| 99 | self.operator_step += 1 |
| 100 | return self |
| 101 | |
| 102 | def reset(self): |
| 103 | self.operator_step = -1 |
| 104 | return self |
| 105 | |
| 106 | def _load_local_file(self, file_path: str, file_type: str) -> pd.DataFrame: |
| 107 | """Load data from local file based on file type.""" |
| 108 | try: |
| 109 | if file_type == "json": |
| 110 | return pd.read_json(file_path) |
| 111 | elif file_type == "jsonl": |
| 112 | import json |
| 113 | records = [] |
| 114 | with open(file_path, 'r', encoding='utf-8') as f: |
| 115 | for line_num, line in enumerate(f, 1): |
| 116 | line = line.strip() |
| 117 | if line: |
| 118 | try: |
| 119 | record = json.loads(line) |
| 120 | records.append(record) |
| 121 | except json.JSONDecodeError as e: |
| 122 | import warnings |
| 123 | warnings.warn( |
| 124 | f"Skipping invalid JSON at line {line_num} in {file_path}: {e}", |
| 125 | UserWarning |
| 126 | ) |
| 127 | continue |
no outgoing calls