校验MyScaleDBStorage实例的关键参数有效性: - pipeline_id, input_task_id, output_task_id 必须非空,否则抛出异常。 - page_size, page_num 若未设置则赋默认值(page_size=10000, page_num=0)。 所有算子在使用storage前应调用本方法。
(self)
| 295 | Storage for Myscale/ClickHouse database using clickhouse_driver. |
| 296 | """ |
| 297 | def validate_required_params(self): |
| 298 | """ |
| 299 | 校验MyScaleDBStorage实例的关键参数有效性: |
| 300 | - pipeline_id, input_task_id, output_task_id 必须非空,否则抛出异常。 |
| 301 | - page_size, page_num 若未设置则赋默认值(page_size=10000, page_num=0)。 |
| 302 | 所有算子在使用storage前应调用本方法。 |
| 303 | """ |
| 304 | missing = [] |
| 305 | if not self.pipeline_id: |
| 306 | missing.append('pipeline_id') |
| 307 | if not self.input_task_id: |
| 308 | missing.append('input_task_id') |
| 309 | if not self.output_task_id: |
| 310 | missing.append('output_task_id') |
| 311 | if missing: |
| 312 | raise ValueError(f"Missing required storage parameters: {', '.join(missing)}") |
| 313 | if not hasattr(self, 'page_size') or self.page_size is None: |
| 314 | self.page_size = 10000 |
| 315 | if not hasattr(self, 'page_num') or self.page_num is None: |
| 316 | self.page_num = 0 |
| 317 | |
| 318 | def __init__( |
| 319 | self, |