feat(data-process): 清理PDF文档级噪声
This commit is contained in:
@@ -37,6 +37,7 @@ from app.modules.data_process.algorithms import (
|
||||
desensitize_pii,
|
||||
desensitize_structured_record,
|
||||
detect_document_structure,
|
||||
detect_pdf_document_noise,
|
||||
estimate_token_count,
|
||||
extract_pdf_page_texts,
|
||||
generate_standard_records,
|
||||
@@ -44,6 +45,7 @@ from app.modules.data_process.algorithms import (
|
||||
near_duplicate_fingerprint,
|
||||
parse_text_content,
|
||||
preprocess_structured_records,
|
||||
remove_document_noise,
|
||||
score_quality,
|
||||
)
|
||||
from app.modules.data_process.generation import generate_model_records
|
||||
@@ -453,13 +455,25 @@ def _build_preview_items(
|
||||
for source in source_files:
|
||||
parsed = _parse_stored_source(source)
|
||||
if process_type == "unstructured":
|
||||
document_noise_spans = tuple(source.get("document_noise_spans") or ())
|
||||
chunks = _chunk_source_text(parsed.text, config, preprocess_options)
|
||||
for chunk, heading_path in chunks:
|
||||
content = (
|
||||
remove_document_noise(
|
||||
chunk.content,
|
||||
document_noise_spans,
|
||||
source_offset=chunk.start,
|
||||
)
|
||||
if should_clean_invalid and document_noise_spans
|
||||
else chunk.content
|
||||
)
|
||||
preprocess_flags = content_quality_flags(
|
||||
chunk.content,
|
||||
content,
|
||||
min_chars=0,
|
||||
min_tokens=0,
|
||||
)
|
||||
if content != chunk.content:
|
||||
preprocess_flags = (*preprocess_flags, "document_noise_removed")
|
||||
flag_set = set(preprocess_flags)
|
||||
if "clean_invalid_content" in preprocess_options and flag_set & {
|
||||
"empty_content",
|
||||
@@ -475,20 +489,19 @@ def _build_preview_items(
|
||||
}:
|
||||
continue
|
||||
if "deduplicate_content" in preprocess_options:
|
||||
band_keys = _near_duplicate_band_keys(chunk.content)
|
||||
band_keys = _near_duplicate_band_keys(content)
|
||||
candidates = {
|
||||
previous
|
||||
for key in band_keys
|
||||
for previous in seen_near_duplicate_bands.get(key, ())
|
||||
}
|
||||
if any(
|
||||
_safe_near_duplicate(chunk.content, previous)
|
||||
_safe_near_duplicate(content, previous)
|
||||
for previous in candidates
|
||||
):
|
||||
continue
|
||||
for key in band_keys:
|
||||
seen_near_duplicate_bands.setdefault(key, []).append(chunk.content)
|
||||
content = chunk.content
|
||||
seen_near_duplicate_bands.setdefault(key, []).append(content)
|
||||
pii_counts: dict[str, int] = {}
|
||||
if should_desensitize:
|
||||
content, pii_counts = desensitize_pii(content)
|
||||
@@ -512,7 +525,7 @@ def _build_preview_items(
|
||||
"source_end": chunk.end,
|
||||
"source_start_line": chunk.start_line,
|
||||
"source_end_line": chunk.end_line,
|
||||
"token_count": chunk.token_count,
|
||||
"token_count": estimate_token_count(content),
|
||||
"status": "modified" if content != chunk.content else "original",
|
||||
"quality_score": quality,
|
||||
}
|
||||
@@ -1234,6 +1247,7 @@ def pull_external_source(
|
||||
def _prepare_preview_items(
|
||||
task_id: str,
|
||||
store: DataProcessStore,
|
||||
storage: LocalDataProcessStorage,
|
||||
source_file_ids: list[str] | None = None,
|
||||
) -> list[dict[str, Any]]:
|
||||
task = store.get_task(task_id)
|
||||
@@ -1251,6 +1265,48 @@ def _prepare_preview_items(
|
||||
]
|
||||
if not sources:
|
||||
raise InvalidStateError("at least one source file is required")
|
||||
preprocess_options = _preprocess_options(task.get("config") or {})
|
||||
if (
|
||||
task.get("process_type") == "unstructured"
|
||||
and preprocess_options & {"clean_invalid", "clean_invalid_content"}
|
||||
):
|
||||
for index, source in enumerate(sources):
|
||||
if str(source.get("file_format") or "").lower() != "pdf":
|
||||
continue
|
||||
storage_object_id = str(source.get("storage_object_id") or "")
|
||||
actual_size = storage.file_size(
|
||||
storage_object_id,
|
||||
expected_task_id=task_id,
|
||||
expected_source_file_id=str(source["id"]),
|
||||
)
|
||||
if actual_size is None:
|
||||
logger.info(
|
||||
"skip PDF document noise detection for unavailable legacy source %s",
|
||||
source["id"],
|
||||
)
|
||||
continue
|
||||
expected_size = int(source.get("size_bytes") or 0)
|
||||
if expected_size and actual_size != expected_size:
|
||||
raise ValueError("source object size does not match metadata")
|
||||
raw = b"".join(
|
||||
storage.iter_bytes(
|
||||
storage_object_id,
|
||||
expected_task_id=task_id,
|
||||
expected_source_file_id=str(source["id"]),
|
||||
expected_size=actual_size,
|
||||
)
|
||||
)
|
||||
pages = extract_pdf_page_texts(raw)
|
||||
extracted_text = "\n\n".join(page.text for page in pages if page.text)
|
||||
if extracted_text != str(source.get("content") or ""):
|
||||
logger.warning(
|
||||
"skip PDF document noise detection because stored offsets differ for %s",
|
||||
source["id"],
|
||||
)
|
||||
continue
|
||||
enriched = dict(source)
|
||||
enriched["document_noise_spans"] = detect_pdf_document_noise(pages)
|
||||
sources[index] = enriched
|
||||
items = _build_preview_items(task, sources)
|
||||
if not items and source_file_ids is None:
|
||||
raise InvalidStateError("source files did not produce preview items")
|
||||
@@ -1262,10 +1318,11 @@ def build_preview(
|
||||
task_id: str,
|
||||
payload: PreviewBuildRequest = Body(default_factory=PreviewBuildRequest),
|
||||
store: DataProcessStore = Depends(get_data_process_store),
|
||||
storage: LocalDataProcessStorage = Depends(get_data_process_storage),
|
||||
) -> dict[str, Any]:
|
||||
with api_errors():
|
||||
selected_ids = payload.source_file_ids
|
||||
items = _prepare_preview_items(task_id, store, selected_ids)
|
||||
items = _prepare_preview_items(task_id, store, storage, selected_ids)
|
||||
created = store.replace_preview_items(
|
||||
task_id,
|
||||
items,
|
||||
@@ -1405,9 +1462,10 @@ def start(
|
||||
background_tasks: BackgroundTasks,
|
||||
payload: GenerateRequest = Body(default_factory=GenerateRequest),
|
||||
store: DataProcessStore = Depends(get_data_process_store),
|
||||
storage: LocalDataProcessStorage = Depends(get_data_process_storage),
|
||||
) -> dict[str, Any]:
|
||||
with api_errors():
|
||||
items = _prepare_preview_items(task_id, store)
|
||||
items = _prepare_preview_items(task_id, store, storage)
|
||||
store.replace_preview_items(task_id, items)
|
||||
return _start_generation(task_id, payload, background_tasks, store)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user