MCPcopy Create free account
hub / github.com/OpenDCAI/DataFlow-MM / read

Method read

dataflow/utils/storage.py:353–387  ·  view source on GitHub ↗

Read data from Myscale/ClickHouse table.

(self, output_type: Literal["dataframe", "dict"])

Source from the content-addressed store, hash-verified

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 """

Callers

nothing calls this directly

Calls

no outgoing calls

Tested by

no test coverage detected