Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions data_server/formatify/FormatifyManager.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,8 @@ def _submit_formatify_task_to_csghub(
namespace: str,
user_name: str | None = None,
task_run_time: str | None = None,
use_streaming: bool | None = None,
chunk_size: int | None = None,
):
flow_id = build_job_flow_id("formatify", formatify_task.id)
formatify_task.flow_id = flow_id
Expand All @@ -73,6 +75,12 @@ def _submit_formatify_task_to_csghub(
task_params["flow_id"] = flow_id
if task_run_time:
task_params["execute_time"] = task_run_time

# Add streaming mode parameters (not stored in database, passed via task_params)
if use_streaming is not None:
task_params["use_streaming"] = use_streaming
if chunk_size is not None:
task_params["chunk_size"] = chunk_size
dag_tasks = build_formatify_dag(flow_id, task_params)
payload = build_csghub_payload(
job_id=flow_id,
Expand Down Expand Up @@ -124,6 +132,10 @@ def create_formatify_task(
# Prepare skip_meta value (use provided value or default to False)
skip_meta_value = dataFormatTask.skip_meta if dataFormatTask.skip_meta is not None else False

# Extract streaming mode parameters (not stored in database)
use_streaming = dataFormatTask.use_streaming
chunk_size = dataFormatTask.chunk_size

data_format_task_db = DataFormatTask(name=dataFormatTask.name,
des=dataFormatTask.des,
from_csg_hub_dataset_name=dataFormatTask.from_csg_hub_dataset_name,
Expand Down Expand Up @@ -163,6 +175,8 @@ def create_formatify_task(
user_token,
nu,
user_name=user_name,
use_streaming=use_streaming,
chunk_size=chunk_size,
)
data_format_task_db.task_status = DataFormatTaskStatusEnum.WAITING.value
except Exception as e:
Expand Down Expand Up @@ -409,6 +423,8 @@ def execute_formatify_task(
user_name: str,
user_token: str,
task_run_time: str | None = None,
use_streaming: bool | None = None,
chunk_size: int | None = None,
):
"""Waiting and not yet submitted to CSGHub: first submit this record."""
try:
Expand All @@ -425,6 +441,8 @@ def execute_formatify_task(
nu,
user_name=user_name,
task_run_time=task_run_time,
use_streaming=use_streaming,
chunk_size=chunk_size,
)
formatify_task.task_status = DataFormatTaskStatusEnum.WAITING.value
db_session.commit()
Expand All @@ -442,6 +460,8 @@ def execute_new_formatify_task(
user_name: str,
user_token: str,
task_run_time: str | None = None,
use_streaming: bool | None = None,
chunk_size: int | None = None,
):
"""List "Execute": copy new task and submit to CSGHub; do not re-run old record."""
try:
Expand All @@ -462,6 +482,8 @@ def execute_new_formatify_task(
nu,
user_name=user_name,
task_run_time=task_run_time,
use_streaming=use_streaming,
chunk_size=chunk_size,
)
new_task.task_status = DataFormatTaskStatusEnum.WAITING.value
db_session.commit()
Expand Down
4 changes: 4 additions & 0 deletions data_server/formatify/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,10 @@ class DataFormatTaskRequest(BaseModel):
storage_size: Optional[str] = None
namespace_uuid: Optional[str] = None
namespace_type: str = "personal"

# Streaming mode parameters (not stored in database, only passed to conversion tasks)
use_streaming: Optional[bool] = None
chunk_size: Optional[int] = None

@field_validator("storage_size", mode="before")
@classmethod
Expand Down
60 changes: 58 additions & 2 deletions data_server/pod/common_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,9 @@
convert_excel_to_csv,
convert_excel_to_json,
convert_excel_to_parquet,
convert_csv_to_json_streaming,
convert_csv_to_parquet_streaming,
convert_csv_to_excel_streaming,
convert_word_to_markdown,
convert_txt_to_markdown,
convert_html_to_markdown,
Expand Down Expand Up @@ -828,7 +831,17 @@ def run_format_conversion(task_params: dict):
if not found:
raise ValueError("No source files found for format conversion")

convert_func = _select_convert_func(task_params.get("from_data_type"), task_params.get("to_data_type"))
# Extract streaming mode parameters
use_streaming = task_params.get("use_streaming", False)
if isinstance(use_streaming, str):
use_streaming = use_streaming.lower() in ("true", "1", "yes")
use_streaming = bool(use_streaming)

convert_func = _select_convert_func(
task_params.get("from_data_type"),
task_params.get("to_data_type"),
use_streaming=use_streaming
)
if convert_func is None:
raise ValueError("Unsupported format conversion")

Expand Down Expand Up @@ -983,7 +996,30 @@ def run_format_conversion(task_params: dict):
}


def _select_convert_func(from_type, to_type):
def _select_convert_func(from_type, to_type, use_streaming=False):
"""
Select conversion function based on format types and mode.

Args:
from_type: Source format type
to_type: Target format type
use_streaming: If True, use streaming mode for CSV conversions (when available)

Returns:
Conversion function or None
"""
# For CSV source with streaming mode enabled, use streaming functions
if use_streaming and from_type == DataFormatTypeEnum.Csv.value:
streaming_mapping = {
DataFormatTypeEnum.Excel.value: convert_csv_to_excel_streaming,
DataFormatTypeEnum.Json.value: convert_csv_to_json_streaming,
DataFormatTypeEnum.Parquet.value: convert_csv_to_parquet_streaming,
}
streaming_func = streaming_mapping.get(to_type)
if streaming_func:
return streaming_func

# Standard (non-streaming) mode mapping
mapping = {
(DataFormatTypeEnum.Excel.value, DataFormatTypeEnum.Csv.value): convert_excel_to_csv,
(DataFormatTypeEnum.Excel.value, DataFormatTypeEnum.Json.value): convert_excel_to_json,
Expand All @@ -1002,13 +1038,33 @@ def _select_convert_func(from_type, to_type):


def _run_convert_func(convert_func, file_path: str, task_uid: str, task_params: dict):

# PDF to Markdown has special parameters
if convert_func is convert_pdf_to_markdown:
return convert_func(
file_path,
task_uid,
task_params.get("mineru_api_url"),
task_params.get("mineru_backend"),
)

# Streaming conversion functions require chunk_size
streaming_funcs = (
convert_csv_to_json_streaming,
convert_csv_to_parquet_streaming,
convert_csv_to_excel_streaming,
)

if convert_func in streaming_funcs:
chunk_size = task_params.get("chunk_size", 50000)
try:
chunk_size = int(chunk_size)
except (ValueError, TypeError):
chunk_size = 50000

return convert_func(file_path, task_uid, chunk_size=chunk_size)

# Standard conversion functions
return convert_func(file_path, task_uid)


Expand Down
Loading