Read data from Myscale/ClickHouse table.
(self, output_type: Literal["dataframe", "dict"])
| 351 | self.validate_required_params() |
| 352 | |
| 353 | def read(self, output_type: Literal["dataframe", "dict"]) -> Any: |
| 354 | """ |
| 355 | Read data from Myscale/ClickHouse table. |
| 356 | """ |
| 357 | where_clauses = [] |
| 358 | params = {} |
| 359 | if self.pipeline_id: |
| 360 | where_clauses.append("pipeline_id = %(pipeline_id)s") |
| 361 | params['pipeline_id'] = self.pipeline_id |
| 362 | if self.input_task_id: |
| 363 | where_clauses.append("task_id = %(task_id)s") |
| 364 | params['task_id'] = self.input_task_id |
| 365 | where_sql = f"WHERE {' AND '.join(where_clauses)}" if where_clauses else "" |
| 366 | limit_offset = f"LIMIT {self.page_size} OFFSET {(self.page_num-1)*self.page_size}" if self.page_size else "" |
| 367 | sql = f"SELECT * FROM {self.table} {where_sql} {limit_offset}" |
| 368 | self.logger.info(f"Reading from DB: {sql} with params {params}") |
| 369 | result = self.client.execute(sql, params, with_column_types=True) |
| 370 | rows, col_types = result |
| 371 | columns = [col[0] for col in col_types] |
| 372 | df = pd.DataFrame(rows, columns=columns) |
| 373 | # 解析 data 字段为 dict |
| 374 | if 'data' not in df.columns: |
| 375 | raise ValueError("Result does not contain required 'data' field.") |
| 376 | |
| 377 | # 只保留 data 字段 |
| 378 | data_series = df['data'].apply(safe_json_loads) |
| 379 | |
| 380 | if output_type == "dataframe": |
| 381 | # 返回只有 data 一列的 DataFrame |
| 382 | return pd.DataFrame({'data': data_series}) |
| 383 | elif output_type == "dict": |
| 384 | # 返回 data 字段的 dict 列表 |
| 385 | return list(data_series) |
| 386 | else: |
| 387 | raise ValueError(f"Unsupported output type: {output_type}") |
| 388 | |
| 389 | def write(self, data: Any) -> Any: |
| 390 | """ |
nothing calls this directly
no outgoing calls
no test coverage detected