首页
/ Docling service_client 实战:用 DoclingServiceClient 调用 docling-serve 完成远程转换、批量处理与切块

Docling service_client 实战:用 DoclingServiceClient 调用 docling-serve 完成远程转换、批量处理与切块

2026-09-06 12:24:05作者:何举烈Damon

本篇基于仓库中 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-remote CLI 使用相同的变量名),来源可以是系统环境变量,也可以是当前工作目录下的 .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-client extra 的 docling-slim

pip install "docling-slim[service-client]"

仓库的 docling-slim 包说明 中同样列出了 service-client extra,其用途即"Docling service client / Remote processing";此外还有一个元包 docling-client,它只是拉取 docling-slim[service-client],等价于安装上式依赖。

从仓库根目录运行示例

四个示例脚本都引用 tests/data/pdf/sources/ 下的样例文档(如 2305.03393v1-pg9.pdfcode_and_formula.pdfpicture_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_statustask_position 队列位置等);
  • watch(timeout):迭代器,逐条产出状态更新直至终态;
  • result(timeout):等待终态后取回结果;任务未结束时先 waitfetch_result
  • 便捷属性 statusqueue_positiondone(终态为 successfailure,见 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 的结果是 RawServiceResultcontent 字节 + 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 + 可选的每项目 optionsheadersmetadata,定义见 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() 的关键区别:

  • 不接受上传文件,只接受源请求对象(内置类型或服务端启用的插件类型);
  • 必须显式提供 targettargettargets 二者只能给一个,否则 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 枚举只有两个值:HYBRIDHIERARCHICALclient.py),对应切块选项模型 HybridChunkerOptions / HierarchicalChunkerOptionschunking 模型);
  • 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=Falsetarget_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.pyWebSocketWatcher / PollingWatcher 及其异步版本)与 job.pyConversionJob

重试与错误映射

_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 并解析配额细节)、TaskTimeoutErrorTaskNotFoundErrorTaskExecutionErrorResultExpiredError/ResultNotReadyErrorArtifactDownloadError 等,便于按类型做针对性处理。

工件下载的安全防护

高层 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)。

验证与参考路径

适用前提与限制:以上行为以当前仓库(docling)中的 docling.service_client 实现为准;docling-serve 服务需已单独部署运行,且部分能力(如插件源、特定 target、存储回退)取决于服务端配置与版本。客户端会声明其 DoclingDocument 版本能力,跨版本响应结构不匹配时抛出 ResponseSchemaMismatchError,升级客户端与服务端时建议保持配套。

登录后查看全文
热门项目推荐
相关项目推荐