Docling service_client 实战:用 DoclingServiceClient 调用 docling-serve 完成远程转换、批量处理与切块
【免费下载链接】doclingGet your documents ready for gen AI项目地址: https://gitcode.com/GitHub_Trending/do/docling
本篇基于仓库中 Client SDK Examples 文档 展开,介绍如何用docling.service_client包中的DoclingServiceClient对接一个已运行的docling-serve实例:从环境变量配置、单文档convert()、并发convert_all(),到任务级submit*API、批量submit_batch()与 RAG 切块chunk()。读完并对照源码后,你可以独立完成远程文档转换、任务生命周期追踪、结果目标(result target)选择与批量源(HTTP/S3/插件)接入,并理解 SDK 内部的自动回退、状态监听与重试机制。
需要强调前提:这些脚本和 SDK本身不启动任何服务,它们假设docling-serve已经在某处运行,客户端只负责发请求、追状态、取结果。
环境准备
配置服务地址与 API Key
客户端与示例脚本读取同一组环境变量(与docling convert-remoteCLI 使用相同的变量名),来源可以是系统环境变量,也可以是当前工作目录下的.env文件:
DOCLING_SERVICE_URL=https://your-docling-service.example.com DOCLING_SERVICE_API_KEY=your-api-key # 服务未启用鉴权时可省略从源码可以看到两者的用途(客户端实现):
url会被_normalize_base_url()严格校验:必须是绝对 http(s) 地址,不允许携带 query/fragment,也不允许把/v1写进 base URL(API 版本号由客户端自己拼接);api_key非空时会被放进每个 HTTP 请求的X-Api-Key请求头,WebSocket 状态订阅时同样以 header 或 query 形式携带;- 客户端还会发送
Accept-Docling-Document-Version头,声明本端 docling-core 能读取的最新DoclingDocument版本,服务端需要时会向下兼容转换(见 client.py 构造函数)。
安装客户端 SDK
按文档要求安装带service-clientextra 的docling-slim:
pip install "docling-slim[service-client]"仓库的 docling-slim 包说明 中同样列出了service-clientextra,其用途即"Docling service client / Remote processing";此外还有一个元包 docling-client,它只是拉取docling-slim[service-client],等价于安装上式依赖。
从仓库根目录运行示例
四个示例脚本都引用tests/data/pdf/sources/下的样例文档(如2305.03393v1-pg9.pdf、code_and_formula.pdf、picture_classification.pdf),这些文件在仓库中真实存在(见 tests/data/pdf/sources),因此文档要求从仓库根目录运行,例如:
uv run python docs/examples/service_client/convert.py示例总览
| 脚本 | 展示内容 |
|---|---|
| convert.py | convert()与convert_all()—— 高层 API |
| tasks.py | submit*系列 API:任务生命周期、结果目标、逐项 fan-out |
| batch.py | submit_batch():内置或插件源与工件(artifact)目标 |
| chunk.py | chunk():把文档切分为检索就绪的片段 |
下面逐一展开,并结合源码说明每个 API 背后的实际行为。
基础 API:convert() 与 convert_all()
单文档转换:convert()
调用形态与本地DocumentConverter一致——传入源,拿回ConversionResult,直接export_to_markdown():
from docling.service_client import DoclingServiceClient client = DoclingServiceClient(url=..., api_key=...) result = client.convert(source="path/to/report.pdf") # 也接受 http(s) URL print(result.document.export_to_markdown())convert.py 的完整写法(含with上下文管理器自动关闭 HTTP 资源、打印前 500 字符 Markdown 预览):
with DoclingServiceClient( url=os.environ["DOCLING_SERVICE_URL"], api_key=os.environ.get("DOCLING_SERVICE_API_KEY", ""), ) as client: result = client.convert(source=SINGLE) print("convert():", result.document.name, result.status.value) print(result.document.export_to_markdown()[:500])默认行为(启用 OCR、表格结构识别、输出 Markdown)与 docling 本地DocumentConverter的默认值一致;只有需要覆盖时才传options=ConvertDocumentsOptions(...)。ConvertDocumentsOptions定义在 服务选项模型,在客户端中被别名为ConvertDocumentsRequestOptions使用。
convert()的完整签名(client.py)还提供了几个实用参数:
headers:附加到提交请求的 HTTP 头;max_num_pages/max_file_size:客户端**预检(preflight)**限制。从源码看,_preflight_limits()会在本地文件超限时直接返回一个SKIPPED状态的ConversionResult(失败类别为POLICY),而不浪费一次网络提交;page_range:页码范围,若与max_num_pages同时给出则取交集;raises_on_error:默认True,当结果为非成功状态(不是SUCCESS/PARTIAL_SUCCESS)时抛出ConversionError。
多文档并发转换:convert_all()
for result in client.convert_all( source=["a.pdf", "b.pdf", "https://.../c.pdf"], max_concurrency=4, ): print(result.input.file.name, result.status)convert_all()返回一个按输入顺序产出的Iterator[ConversionResult]。要点(对照 convert_all 实现):
max_concurrency控制同时提交的任务数;不传时使用客户端构造参数max_concurrency,其默认值为 8(DEFAULT_MAX_CONCURRENCY),上限 512(MAX_CONCURRENCY_LIMIT),越界会直接ValueError;- 旧参数名
sources=已弃用(DeprecationWarning),请使用source=; - 单个文档的失败不会中断整个迭代:从
_convert_all_async()可以看到,失败的项会被降级为一个FAILURE状态的ConversionResult并继续产出后续结果,因此批量作业中应按result.status逐项判断; - 同步版
convert_all()内部通过私有事件循环驱动原生异步客户端(_build_async_service_client()),即并发 fan-out 不依赖线程池,且不能在已运行的 asyncio loop 中调用(会抛RuntimeError,见_ensure_sync_bridge_allowed())。
支持的源类型
从类型别名SourceType: Path | str | DocumentStream | HttpSourceRequest看,convert()/convert_all()/chunk()接受:本地Path、http(s) URL 字符串(会被规范化为HttpSourceRequest,非 http/https 方案会报Unsupported URL scheme)、内存中的DocumentStream,或显式的HttpSourceRequest(可携带请求头)。提交时本地/流式源走multipart文件上传端点/v1/convert/file/async,HTTP 源走 JSON 端点/v1/convert/source/async(见 _submit_convert_task)。
submit* 系列:任务生命周期与结果目标(tasks.py)
convert()/convert_all()返回的是"重建好"的ConversionResult;而submit*家族返回原始服务响应,并让你显式决定结果放在哪里。tasks.py 覆盖三种用法。
任务生命周期:submit() -> watch() -> result()
job = client.submit(source=SOURCE, output_formats=[OutputFormat.MARKDOWN]) print("task id:", job.task_id) for update in job.watch(timeout=300.0): print(" status:", update.task_status, "position:", update.task_position) result = job.result(timeout=300.0) print("done:", result.num_succeeded, "succeeded /", result.num_failed, "failed")submit()立即返回一个ConversionJob句柄,定义在 job.py。它暴露:
task_id/submitted_at:任务标识与提交时间;poll(wait):单次拉取TaskStatusResponse(含task_status、task_position队列位置等);watch(timeout):迭代器,逐条产出状态更新直至终态;result(timeout):等待终态后取回结果;任务未结束时先wait再fetch_result。- 便捷属性
status、queue_position、done(终态为success或failure,见 watchers.py 中 TERMINAL_TASK_STATUSES)。
关于"省略target"的行为:submit()在未指定 target 时默认使用PresignedUrlTarget()(见 submit 实现 中resolved_target = PresignedUrlTarget() if target is None else target);对高层convert()/convert_all()而言则是"presigned 优先、失败时回退 in-body"的自动目标(下一节详述)。
结果目标(result targets)
# ZipTarget:返回请求格式产物的原始压缩包 archive = client.submit( source=SOURCE, output_formats=[OutputFormat.MARKDOWN], target=ZipTarget(), ).result(timeout=300.0) print("ZipTarget:", archive.content_type, len(archive.content), "bytes")ZipTarget的结果是RawServiceResult(content字节 +content_type+ 文件名)。SDK 支持的目标类型(client.py 顶部 TypeAlias):
| 目标 | 结果形态 |
|---|---|
InBodyTarget | 文档内容直接内联在响应体中(ConvertDocumentResponse);tasks.py 注释指出部分服务会限制目标类型,仅接受存储后端类型 |
PresignedUrlTarget | 返回下载 URL(PresignedUrlConvertResponse) |
ZipTarget | 返回原始 ZIP 压缩包(RawServiceResult) |
S3Target/AzureBlobTarget/GoogleCloudStorageTarget/GoogleDriveTarget | 存储类目标,结果写入外部存储,返回PresignedUrlConvertDocumentResponse |
一个容易忽略的细节:当目标是InBodyTarget(或客户端要重建ConversionResult)时,_options_for_output_formats()会自动在输出格式中补上OutputFormat.JSON——因为重建DoclingDocument需要 JSON 文档载荷(见 _with_json_output_format)。
逐项 fan-out:submit_and_retrieve_each()
items = [ ConversionItem(source=s, metadata={"id": i}) for i, s in enumerate(MANY) ] for item, outcome in client.submit_and_retrieve_each(items, max_in_flight=4): if isinstance(outcome, Exception): print(" ", item.metadata, "failed:", outcome) else: print(" ", item.metadata, "ok")submit_and_retrieve_each()接收ConversionItem列表(source+ 可选的每项目options、headers、metadata,定义见 client.py),以max_in_flight并发度提交,每个输入对应一个结果:成功时是相应的响应模型,失败时异常本身作为 outcome 内联产出(而不是抛出中断)。target=None时采用与submit()相同的自动目标策略;ordered参数控制产出顺序。旧方法名submit_and_retrieve_many()已弃用并转发到submit_and_retrieve_each()。从实现看(_submit_and_retrieve_many_uses_websocket_wait()),当使用 WebSocket 监听且max_in_flight <= 64时,fan-out 会复用 WebSocket 状态流等待,减少轮询开销。
批量转换:submit_batch()(batch.py)
批量端点面向"大量或长时运行"的源,与submit()的关键区别:
- 不接受上传文件,只接受源请求对象(内置类型或服务端启用的插件类型);
- 必须显式提供 target(
target与targets二者只能给一个,否则ValueError,见 submit_batch),提交到POST /v1/convert/source/batch。
batch.py 用两个 arXiv 论文 URL 演示最简流程:
job = client.submit_batch( sources=[AnyHttpSourceRequest(url=url) for url in SOURCES], target=PresignedUrlTarget(), output_formats=[OutputFormat.MARKDOWN, OutputFormat.JSON], ) result = job.result(timeout=300.0) for document in result.documents: print(document.filename, document.status.value) for artifact in document.artifacts: print(" ", artifact.artifact_type, str(artifact.uri))脚本中还给出了两段注释掉的进阶用法,值得留意:
- S3 扇出:
S3SourceRequest读桶内输入、S3Target写回结果桶(需要真实凭证,故注释保留); - 插件连接器:当 SDK 不认识插件 schema 时,
sources/target可以直接传原始 dict mapping(如{"kind": "filenet", ...}),服务端负责完整校验,并且插件连接器必须在服务端显式启用。
切块:chunk()(chunk.py)
对 RAG 场景,SDK 把"转换 + 切分"合并成一次调用:
from docling.service_client import ChunkerKind, DoclingServiceClient response = client.chunk(source=SOURCE, chunker=ChunkerKind.HIERARCHICAL) print(len(response.chunks), "chunks from", len(response.documents), "document(s)") for chunk in response.chunks[:3]: print("---") print(chunk.text[:300])ChunkerKind枚举只有两个值:HYBRID与HIERARCHICAL(client.py),对应切块选项模型HybridChunkerOptions/HierarchicalChunkerOptions(chunking 模型);chunk()本身是便捷封装:submit_chunk()提交任务后直接job.result(timeout=self._job_timeout)取回ChunkDocumentResponse,其中包含chunks与来源documents;- 从
_submit_chunk_task()可以看到底层端点为/v1/chunk/{hybrid|hierarchical}/file/async(文件上传)或/v1/chunk/{...}/source/async(HTTP 源),且固定include_converted_doc=False、target_type=inbody——即切块作业默认不回传完整转换文档。
源码级机制:自动目标回退、状态监听与容错
以下行为都直接影响你在生产环境中的可观测性与稳定性,均来自 docling/service_client 包的实现。
自动目标回退(presigned -> in-body)
高层convert()/convert_all()内部先以PresignedUrlTarget()提交;若服务端因"未配置 artifact 存储"等原因拒绝(400/422 且 detail 含相应提示),_should_fallback_from_presigned_target()判定后自动改用InBodyTarget()重新提交(见 _submit_conversion_job_with_auto_target)。这就是示例中"省略 target"能同时兼容不同服务部署配置的原因。取回结果时,若走的是 presigned 路径,客户端会下载工件并重建ConversionResult:优先选择resource_bundle(ZIP 资源包,解压后内联图片,且带 zip-slip 与越界引用防护),其次退而选择自包含的 JSON 工件(见_select_artifact()/_reconstruct_document_from_bundle())。
状态监听:WebSocket 优先,轮询兜底
DoclingServiceClient构造时status_watcher默认为StatusWatcherKind.WEBSOCKET:
- WebSocket 监听连接
WS /v1/status/ws/{task_id},对连接断开有最多 3 次指数退避重连; - 若 WebSocket 不可用且
ws_fallback_to_poll=True(默认),自动切到轮询监听GET /v1/status/poll/{task_id}?wait=...(服务端支持带等待的长轮询,poll_server_wait默认 5 秒); - 也可以显式传
status_watcher=StatusWatcherKind.POLLING完全走轮询。
相关实现分布在 watchers.py(WebSocketWatcher/PollingWatcher及其异步版本)与 job.py 的ConversionJob。
重试与错误映射
_request_with_retry()对每个 HTTP 调用执行统一重试策略:500/502 走指数退避(基准 1 秒,2**attempt),429/503 优先读取Retry-After响应头;仅GET/HEAD/OPTIONS的传输层错误会重试。默认http_retries=3。错误会被映射为具体异常类型(exceptions 模块):ServiceError(4xx)、ServiceUnavailableError(5xx)、UsageLimitExceededError(402 并解析配额细节)、TaskTimeoutError、TaskNotFoundError、TaskExecutionError、ResultExpiredError/ResultNotReadyError、ArtifactDownloadError等,便于按类型做针对性处理。
工件下载的安全防护
高层 API 下载 presigned 工件时(_download_artifact_bytes()):使用独立 httpx 客户端(不携带服务X-Api-Key)、手动跟踪重定向且每一跳都过 SSRF 校验(拒绝内网/回环等地址)、流式累计字节数超过max_artifact_download_bytes(默认 512 MiB)即中止、重定向超过 5 次报错。对私有/内网存储场景存在内部开关_allow_private_artifact_urls(默认关闭)。
异步客户端
同包提供AsyncDoclingServiceClient(_async_client.py)与AsyncConversionJob,API 形态与同步版一一对应(async for/await)。同步版的批量接口实际上就是在私有事件循环中驱动它完成并发,二者共享同一套 URL 校验、选项序列化与重试逻辑(_BaseDoclingServiceClient)。
验证与参考路径
- 示例文档:docs/examples/service_client/README.md 及 convert.py、tasks.py、batch.py、chunk.py;
- SDK 实现:docling/service_client/client.py、job.py、watchers.py、exceptions.py、_async_client.py;
- 请求/响应/目标数据模型:docling/datamodel/service/(
requests.py、responses.py、targets.py、options.py、chunking.py); - 测试:SDK 单元测试 tests/test_service_client_sdk_unit.py、基于模拟服务的集成测试 tests/test_service_client_fake_service.py(模拟服务见 tests/fakes/docling_serve.py),另有 test_service_client_payload_fidelity.py 等;
- 样例输入文档:tests/data/pdf/sources/。
适用前提与限制:以上行为以当前仓库(docling)中的docling.service_client实现为准;docling-serve服务需已单独部署运行,且部分能力(如插件源、特定 target、存储回退)取决于服务端配置与版本。客户端会声明其DoclingDocument版本能力,跨版本响应结构不匹配时抛出ResponseSchemaMismatchError,升级客户端与服务端时建议保持配套。
【免费下载链接】doclingGet your documents ready for gen AI项目地址: https://gitcode.com/GitHub_Trending/do/docling
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考