mirror of
https://github.com/Tencent/WeKnora.git
synced 2026-08-30 16:53:21 +08:00
fix(docreader): throttle heavy parser concurrency
This commit is contained in:
+8
-1
@@ -398,6 +398,13 @@ DOCREADER_TRANSPORT=grpc
|
||||
# 设为正整数(如 500)可限制超大 Word 文档的解析开销;超过页数的内容将不会继续解析
|
||||
# DOCREADER_DOCX_MAX_PAGES=0
|
||||
|
||||
# Docreader 重型解析器并发控制,适合限制多 PDF/OCR 同时导入时的资源占用
|
||||
# 设为 0 可关闭对应限流;默认 1 更适合普通单机或 GPU 资源有限的部署
|
||||
# DOCREADER_MARKITDOWN_MAX_WORKERS=1
|
||||
# DOCREADER_PDF_RENDER_MAX_WORKERS=1
|
||||
# DOCREADER_PDF_RENDER_DPI=200
|
||||
# DOCREADER_PDF_JPEG_QUALITY=90
|
||||
|
||||
# 如果使用Weaviate作为向量存储,需要配置以下参数
|
||||
# 注意:容器内访问请使用 service:port(不要用 localhost,也不要用宿主机映射端口)
|
||||
# Weaviate HTTP 地址(Docker 内:weaviate:8080;宿主机访问:localhost:9035)
|
||||
@@ -442,4 +449,4 @@ DOCREADER_TRANSPORT=grpc
|
||||
|
||||
# 用于OIDC用于信息中提取用户数据
|
||||
# OIDC_USER_INFO_MAPPING_USER_NAME=name
|
||||
# OIDC_USER_INFO_MAPPING_EMAIL=email
|
||||
# OIDC_USER_INFO_MAPPING_EMAIL=email
|
||||
|
||||
@@ -149,6 +149,10 @@ services:
|
||||
- DOCREADER_IMAGE_OUTPUT_DIR=/tmp/docreader
|
||||
- MINERU_ENDPOINT=${MINERU_ENDPOINT:-}
|
||||
- MAX_FILE_SIZE_MB=${MAX_FILE_SIZE_MB:-}
|
||||
- DOCREADER_MARKITDOWN_MAX_WORKERS=${DOCREADER_MARKITDOWN_MAX_WORKERS:-1}
|
||||
- DOCREADER_PDF_RENDER_MAX_WORKERS=${DOCREADER_PDF_RENDER_MAX_WORKERS:-1}
|
||||
- DOCREADER_PDF_RENDER_DPI=${DOCREADER_PDF_RENDER_DPI:-200}
|
||||
- DOCREADER_PDF_JPEG_QUALITY=${DOCREADER_PDF_JPEG_QUALITY:-90}
|
||||
healthcheck:
|
||||
test: ["CMD", "grpc_health_probe", "-addr=localhost:50051"]
|
||||
interval: 30s
|
||||
|
||||
@@ -185,6 +185,10 @@ services:
|
||||
environment:
|
||||
- DOCREADER_IMAGE_OUTPUT_DIR=/tmp/docreader
|
||||
- MAX_FILE_SIZE_MB=${MAX_FILE_SIZE_MB:-}
|
||||
- DOCREADER_MARKITDOWN_MAX_WORKERS=${DOCREADER_MARKITDOWN_MAX_WORKERS:-1}
|
||||
- DOCREADER_PDF_RENDER_MAX_WORKERS=${DOCREADER_PDF_RENDER_MAX_WORKERS:-1}
|
||||
- DOCREADER_PDF_RENDER_DPI=${DOCREADER_PDF_RENDER_DPI:-200}
|
||||
- DOCREADER_PDF_JPEG_QUALITY=${DOCREADER_PDF_JPEG_QUALITY:-90}
|
||||
healthcheck:
|
||||
test: ["CMD", "grpc_health_probe", "-addr=localhost:50051"]
|
||||
interval: 30s
|
||||
|
||||
+9
-1
@@ -77,6 +77,13 @@ docreader:
|
||||
- `DOCREADER_GRPC_MAX_WORKERS`: gRPC 服务的最大工作线程数(默认:4)
|
||||
- `DOCREADER_GRPC_PORT`: gRPC 服务监听端口(默认:50051)
|
||||
|
||||
### 解析器资源控制
|
||||
|
||||
- `DOCREADER_MARKITDOWN_MAX_WORKERS`: MarkItDown 解析的最大并发数(默认:1,设为 0 可关闭限流)
|
||||
- `DOCREADER_PDF_RENDER_MAX_WORKERS`: 扫描 PDF 渲染为图片的最大并发数(默认:1,设为 0 可关闭限流)
|
||||
- `DOCREADER_PDF_RENDER_DPI`: 扫描 PDF 渲染 DPI(默认:200)
|
||||
- `DOCREADER_PDF_JPEG_QUALITY`: 扫描 PDF 输出 JPEG 质量(默认:90,范围会自动限制在 1-95)
|
||||
|
||||
### OCR 配置
|
||||
|
||||
- `OCR_BACKEND`: OCR 引擎后端,可选值:
|
||||
@@ -145,7 +152,8 @@ DocReader 支持多种存储后端:
|
||||
|
||||
### 图像处理配置
|
||||
|
||||
- `IMAGE_MAX_CONCURRENT`: 图像处理的最大并发数(默认:1)
|
||||
扫描 PDF 会被渲染为 JPEG 图片后交给 Go App 侧 OCR 处理。如果在导入多个大 PDF 时出现资源占用过高,
|
||||
可以优先调低 `DOCREADER_PDF_RENDER_MAX_WORKERS` 或 `DOCREADER_MARKITDOWN_MAX_WORKERS`。
|
||||
|
||||
## 配置示例
|
||||
|
||||
|
||||
@@ -54,6 +54,10 @@ class DocReaderConfig:
|
||||
|
||||
# Parser
|
||||
docx_max_pages: int
|
||||
markitdown_max_workers: int
|
||||
pdf_render_max_workers: int
|
||||
pdf_render_dpi: int
|
||||
pdf_jpeg_quality: int
|
||||
|
||||
# Proxy
|
||||
external_http_proxy: str
|
||||
@@ -74,6 +78,10 @@ def load_config() -> DocReaderConfig:
|
||||
)
|
||||
grpc_port = _get_int(["DOCREADER_GRPC_PORT", "PORT"], 50051)
|
||||
docx_max_pages = _get_int(["DOCREADER_DOCX_MAX_PAGES"], 0)
|
||||
markitdown_max_workers = _get_int(["DOCREADER_MARKITDOWN_MAX_WORKERS"], 1)
|
||||
pdf_render_max_workers = _get_int(["DOCREADER_PDF_RENDER_MAX_WORKERS"], 1)
|
||||
pdf_render_dpi = _get_int(["DOCREADER_PDF_RENDER_DPI"], 200)
|
||||
pdf_jpeg_quality = _get_int(["DOCREADER_PDF_JPEG_QUALITY"], 90)
|
||||
|
||||
external_http_proxy = _get_str(
|
||||
["DOCREADER_EXTERNAL_HTTP_PROXY", "EXTERNAL_HTTP_PROXY"], ""
|
||||
@@ -91,6 +99,10 @@ def load_config() -> DocReaderConfig:
|
||||
grpc_max_file_size_mb=grpc_max_file_size_mb,
|
||||
grpc_port=grpc_port,
|
||||
docx_max_pages=docx_max_pages,
|
||||
markitdown_max_workers=markitdown_max_workers,
|
||||
pdf_render_max_workers=pdf_render_max_workers,
|
||||
pdf_render_dpi=pdf_render_dpi,
|
||||
pdf_jpeg_quality=pdf_jpeg_quality,
|
||||
external_http_proxy=external_http_proxy,
|
||||
external_https_proxy=external_https_proxy,
|
||||
image_output_dir=image_output_dir,
|
||||
@@ -107,6 +119,10 @@ def dump_config(mask_secrets: bool = True) -> Dict[str, Any]:
|
||||
"DOCREADER_GRPC_MAX_FILE_SIZE_MB": cfg.grpc_max_file_size_mb,
|
||||
"DOCREADER_GRPC_PORT": cfg.grpc_port,
|
||||
"DOCREADER_DOCX_MAX_PAGES": cfg.docx_max_pages,
|
||||
"DOCREADER_MARKITDOWN_MAX_WORKERS": cfg.markitdown_max_workers,
|
||||
"DOCREADER_PDF_RENDER_MAX_WORKERS": cfg.pdf_render_max_workers,
|
||||
"DOCREADER_PDF_RENDER_DPI": cfg.pdf_render_dpi,
|
||||
"DOCREADER_PDF_JPEG_QUALITY": cfg.pdf_jpeg_quality,
|
||||
"DOCREADER_EXTERNAL_HTTP_PROXY": cfg.external_http_proxy,
|
||||
"DOCREADER_EXTERNAL_HTTPS_PROXY": cfg.external_https_proxy,
|
||||
"DOCREADER_IMAGE_OUTPUT_DIR": cfg.image_output_dir,
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
import logging
|
||||
import threading
|
||||
from contextlib import contextmanager
|
||||
from typing import Dict, Iterator
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_LIMITERS: Dict[str, threading.BoundedSemaphore] = {}
|
||||
_LIMITERS_LOCK = threading.Lock()
|
||||
|
||||
|
||||
def _get_limiter(name: str, max_workers: int) -> threading.BoundedSemaphore:
|
||||
with _LIMITERS_LOCK:
|
||||
limiter = _LIMITERS.get(name)
|
||||
if limiter is None:
|
||||
limiter = threading.BoundedSemaphore(max_workers)
|
||||
_LIMITERS[name] = limiter
|
||||
return limiter
|
||||
|
||||
|
||||
@contextmanager
|
||||
def parser_worker_limit(name: str, max_workers: int) -> Iterator[None]:
|
||||
"""Limit concurrent access to heavy, process-wide parser operations.
|
||||
|
||||
Set max_workers <= 0 to disable throttling for deployments that have enough
|
||||
CPU/GPU resources and know the parser backend is safe under concurrency.
|
||||
"""
|
||||
|
||||
if max_workers <= 0:
|
||||
yield
|
||||
return
|
||||
|
||||
limiter = _get_limiter(name, max_workers)
|
||||
logger.debug("Waiting for %s parser slot (max_workers=%d)", name, max_workers)
|
||||
limiter.acquire()
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
limiter.release()
|
||||
@@ -3,9 +3,11 @@ import logging
|
||||
|
||||
from markitdown import MarkItDown
|
||||
|
||||
from docreader.config import CONFIG
|
||||
from docreader.models.document import Document
|
||||
from docreader.parser.base_parser import BaseParser
|
||||
from docreader.parser.chain_parser import PipelineParser
|
||||
from docreader.parser.concurrency import parser_worker_limit
|
||||
from docreader.parser.markdown_parser import MarkdownParser
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -33,14 +35,14 @@ class StdMarkitdownParser(BaseParser):
|
||||
if ext and not ext.startswith('.'):
|
||||
ext = '.' + ext
|
||||
|
||||
# 直接调用 convert,移除 try-catch,让异常由上层 PipelineParser 统一捕获
|
||||
result = self.markitdown.convert(
|
||||
io.BytesIO(content),
|
||||
file_extension=ext,
|
||||
keep_data_uris=True
|
||||
)
|
||||
with parser_worker_limit("markitdown", CONFIG.markitdown_max_workers):
|
||||
result = self.markitdown.convert(
|
||||
io.BytesIO(content),
|
||||
file_extension=ext,
|
||||
keep_data_uris=True
|
||||
)
|
||||
return Document(content=result.text_content)
|
||||
|
||||
|
||||
class MarkitdownParser(PipelineParser):
|
||||
_parser_cls = (StdMarkitdownParser, MarkdownParser)
|
||||
_parser_cls = (StdMarkitdownParser, MarkdownParser)
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
from docreader.models.document import Document
|
||||
from docreader.config import CONFIG
|
||||
from docreader.parser.base_parser import BaseParser
|
||||
from docreader.parser.chain_parser import FirstParser
|
||||
from docreader.parser.concurrency import parser_worker_limit
|
||||
from docreader.parser.markitdown_parser import MarkitdownParser
|
||||
|
||||
import io
|
||||
@@ -10,36 +12,74 @@ import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _close_pdfium_resource(resource) -> None:
|
||||
close = getattr(resource, "close", None)
|
||||
if close:
|
||||
close()
|
||||
|
||||
|
||||
def _normalize_image_quality(quality: int) -> int:
|
||||
return min(95, max(1, quality))
|
||||
|
||||
|
||||
class PDFScannedParser(BaseParser):
|
||||
"""Fallback parser for scanned PDFs.
|
||||
|
||||
If the primary parser extracts no text (e.g. Markitdown on a scanned PDF),
|
||||
this parser converts each page into an image. The Go App will then perform
|
||||
OCR on the extracted images.
|
||||
this parser converts each page into a JPEG image. The Go App will then
|
||||
perform OCR on the extracted images.
|
||||
"""
|
||||
|
||||
def parse_into_text(self, content: bytes) -> Document:
|
||||
import pdfplumber
|
||||
import pypdfium2 as pdfium
|
||||
|
||||
images = {}
|
||||
markdown_lines = []
|
||||
|
||||
base_name = os.path.splitext(self.file_name or "document")[0]
|
||||
|
||||
logger.info("PDFScannedParser: Attempting to convert PDF pages to images for %s", self.file_name)
|
||||
logger.info(
|
||||
"PDFScannedParser: Rendering PDF pages to JPEG images for %s",
|
||||
self.file_name,
|
||||
)
|
||||
|
||||
try:
|
||||
with pdfplumber.open(io.BytesIO(content)) as pdf:
|
||||
for i, page in enumerate(pdf.pages):
|
||||
img_obj = page.to_image(resolution=150).original
|
||||
img_byte_arr = io.BytesIO()
|
||||
img_obj.save(img_byte_arr, format="PNG")
|
||||
img_bytes = img_byte_arr.getvalue()
|
||||
|
||||
page_filename = f"{base_name}_page_{i+1}.png"
|
||||
ref_path = f"images/{page_filename}"
|
||||
|
||||
markdown_lines.append(f"")
|
||||
images[ref_path] = base64.b64encode(img_bytes).decode("utf-8")
|
||||
with parser_worker_limit("pdf_render", CONFIG.pdf_render_max_workers):
|
||||
pdf = pdfium.PdfDocument(content)
|
||||
try:
|
||||
page_count = len(pdf)
|
||||
scale = max(1, CONFIG.pdf_render_dpi) / 72
|
||||
quality = _normalize_image_quality(CONFIG.pdf_jpeg_quality)
|
||||
|
||||
for i in range(page_count):
|
||||
page = pdf[i]
|
||||
bitmap = None
|
||||
try:
|
||||
bitmap = page.render(scale=scale)
|
||||
img_obj = bitmap.to_pil()
|
||||
if img_obj.mode != "RGB":
|
||||
img_obj = img_obj.convert("RGB")
|
||||
|
||||
img_byte_arr = io.BytesIO()
|
||||
img_obj.save(
|
||||
img_byte_arr,
|
||||
format="JPEG",
|
||||
quality=quality,
|
||||
optimize=True,
|
||||
)
|
||||
img_bytes = img_byte_arr.getvalue()
|
||||
finally:
|
||||
_close_pdfium_resource(bitmap)
|
||||
_close_pdfium_resource(page)
|
||||
|
||||
page_filename = f"{base_name}_page_{i+1}.jpg"
|
||||
ref_path = f"images/{page_filename}"
|
||||
|
||||
markdown_lines.append(f"")
|
||||
images[ref_path] = base64.b64encode(img_bytes).decode("utf-8")
|
||||
finally:
|
||||
_close_pdfium_resource(pdf)
|
||||
|
||||
text = "\n\n".join(markdown_lines)
|
||||
return Document(
|
||||
@@ -47,11 +87,11 @@ class PDFScannedParser(BaseParser):
|
||||
images=images,
|
||||
metadata={
|
||||
"image_source_type": "scanned_pdf",
|
||||
"page_count": len(pdf.pages)
|
||||
"page_count": page_count
|
||||
}
|
||||
)
|
||||
except Exception as e:
|
||||
logger.exception("PDFScannedParser failed to parse PDF: %v", e)
|
||||
logger.exception("PDFScannedParser failed to parse PDF: %s", e)
|
||||
raise e
|
||||
|
||||
class PDFParser(FirstParser):
|
||||
|
||||
@@ -31,6 +31,7 @@ dependencies = [
|
||||
"pydantic>=2.12.3",
|
||||
"pypdf>=6.1.3",
|
||||
"pypdf2>=3.0.1",
|
||||
"pypdfium2>=5.0.0",
|
||||
"python-docx>=1.2.0",
|
||||
"requests>=2.32.5",
|
||||
"textract==1.5.0",
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
import os
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
|
||||
from docreader import config
|
||||
|
||||
|
||||
class DocReaderConfigTest(unittest.TestCase):
|
||||
def test_parser_concurrency_defaults_are_conservative(self):
|
||||
with patch.dict(os.environ, {}, clear=True):
|
||||
cfg = config.load_config()
|
||||
|
||||
self.assertEqual(cfg.markitdown_max_workers, 1)
|
||||
self.assertEqual(cfg.pdf_render_max_workers, 1)
|
||||
self.assertEqual(cfg.pdf_render_dpi, 200)
|
||||
self.assertEqual(cfg.pdf_jpeg_quality, 90)
|
||||
|
||||
def test_loads_parser_concurrency_env(self):
|
||||
env = {
|
||||
"DOCREADER_MARKITDOWN_MAX_WORKERS": "3",
|
||||
"DOCREADER_PDF_RENDER_MAX_WORKERS": "2",
|
||||
"DOCREADER_PDF_RENDER_DPI": "180",
|
||||
"DOCREADER_PDF_JPEG_QUALITY": "85",
|
||||
}
|
||||
with patch.dict(os.environ, env):
|
||||
cfg = config.load_config()
|
||||
|
||||
self.assertEqual(cfg.markitdown_max_workers, 3)
|
||||
self.assertEqual(cfg.pdf_render_max_workers, 2)
|
||||
self.assertEqual(cfg.pdf_render_dpi, 180)
|
||||
self.assertEqual(cfg.pdf_jpeg_quality, 85)
|
||||
|
||||
def test_dump_config_includes_parser_limits(self):
|
||||
dumped = config.dump_config()
|
||||
|
||||
self.assertIn("DOCREADER_MARKITDOWN_MAX_WORKERS", dumped)
|
||||
self.assertIn("DOCREADER_PDF_RENDER_MAX_WORKERS", dumped)
|
||||
self.assertIn("DOCREADER_PDF_RENDER_DPI", dumped)
|
||||
self.assertIn("DOCREADER_PDF_JPEG_QUALITY", dumped)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,88 @@
|
||||
import base64
|
||||
import io
|
||||
import threading
|
||||
import time
|
||||
import unittest
|
||||
import uuid
|
||||
|
||||
from PIL import Image
|
||||
|
||||
from docreader.parser.concurrency import parser_worker_limit
|
||||
from docreader.parser.pdf_parser import PDFScannedParser, _normalize_image_quality
|
||||
|
||||
|
||||
class ParserConcurrencyTest(unittest.TestCase):
|
||||
def test_parser_worker_limit_serializes_work(self):
|
||||
limiter_name = f"test-{uuid.uuid4()}"
|
||||
active_workers = 0
|
||||
max_active_workers = 0
|
||||
state_lock = threading.Lock()
|
||||
start = threading.Barrier(3)
|
||||
|
||||
def worker():
|
||||
nonlocal active_workers, max_active_workers
|
||||
start.wait()
|
||||
with parser_worker_limit(limiter_name, 1):
|
||||
with state_lock:
|
||||
active_workers += 1
|
||||
max_active_workers = max(max_active_workers, active_workers)
|
||||
time.sleep(0.02)
|
||||
with state_lock:
|
||||
active_workers -= 1
|
||||
|
||||
threads = [threading.Thread(target=worker) for _ in range(2)]
|
||||
for thread in threads:
|
||||
thread.start()
|
||||
|
||||
start.wait()
|
||||
for thread in threads:
|
||||
thread.join()
|
||||
|
||||
self.assertEqual(max_active_workers, 1)
|
||||
|
||||
def test_scanned_pdf_parser_outputs_jpeg_images(self):
|
||||
pdf_bytes = io.BytesIO()
|
||||
pages = [
|
||||
Image.new("RGB", (64, 64), "white"),
|
||||
Image.new("RGB", (64, 64), "black"),
|
||||
]
|
||||
pages[0].save(
|
||||
pdf_bytes,
|
||||
format="PDF",
|
||||
save_all=True,
|
||||
append_images=pages[1:],
|
||||
)
|
||||
|
||||
document = PDFScannedParser(file_name="scan.pdf").parse_into_text(
|
||||
pdf_bytes.getvalue()
|
||||
)
|
||||
|
||||
image_ref = "images/scan_page_1.jpg"
|
||||
self.assertIn(f"", document.content)
|
||||
self.assertIn(image_ref, document.images)
|
||||
self.assertEqual(document.metadata["image_source_type"], "scanned_pdf")
|
||||
self.assertEqual(document.metadata["page_count"], 2)
|
||||
self.assertEqual(len(document.images), 2)
|
||||
self.assertIn("images/scan_page_2.jpg", document.images)
|
||||
image_bytes = base64.b64decode(document.images[image_ref])
|
||||
self.assertTrue(image_bytes.startswith(b"\xff\xd8"))
|
||||
|
||||
def test_scanned_pdf_parser_logs_malformed_pdf_without_format_error(self):
|
||||
parser = PDFScannedParser(file_name="broken.pdf")
|
||||
|
||||
with self.assertLogs("docreader.parser.pdf_parser", level="ERROR") as logs:
|
||||
with self.assertRaises(Exception):
|
||||
parser.parse_into_text(b"not a pdf")
|
||||
|
||||
self.assertTrue(
|
||||
any("PDFScannedParser failed to parse PDF:" in line for line in logs.output)
|
||||
)
|
||||
|
||||
def test_normalize_image_quality_bounds_jpeg_quality(self):
|
||||
self.assertEqual(_normalize_image_quality(-1), 1)
|
||||
self.assertEqual(_normalize_image_quality(90), 90)
|
||||
self.assertEqual(_normalize_image_quality(120), 95)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
Generated
+2
@@ -722,6 +722,7 @@ dependencies = [
|
||||
{ name = "pydantic" },
|
||||
{ name = "pypdf" },
|
||||
{ name = "pypdf2" },
|
||||
{ name = "pypdfium2" },
|
||||
{ name = "python-docx" },
|
||||
{ name = "requests" },
|
||||
{ name = "textract" },
|
||||
@@ -757,6 +758,7 @@ requires-dist = [
|
||||
{ name = "pydantic", specifier = ">=2.12.3" },
|
||||
{ name = "pypdf", specifier = ">=6.1.3" },
|
||||
{ name = "pypdf2", specifier = ">=3.0.1" },
|
||||
{ name = "pypdfium2", specifier = ">=5.0.0" },
|
||||
{ name = "python-docx", specifier = ">=1.2.0" },
|
||||
{ name = "requests", specifier = ">=2.32.5" },
|
||||
{ name = "textract", specifier = "==1.5.0" },
|
||||
|
||||
Reference in New Issue
Block a user