尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Unstract Executor Worker 深度解析:Celery 驱动的 LLM 提取、索引与 Prompt 执行引擎

Unstract Executor Worker 深度解析:Celery 驱动的 LLM 提取、索引与 Prompt 执行引擎 Unstract Executor Worker 深度解析Celery 驱动的 LLM 提取、索引与 Prompt 执行引擎【免费下载链接】unstractLLM-Driven Extraction of Unstructured Data — Built for API Deployments ETL Pipeline Workflows项目地址: https://gitcode.com/GitHub_Trending/un/unstract本篇技术指南聚焦 Unstract 仓库中的 Executor Worker执行器工作进程讲解它如何作为 Celery 工作进程承接 Prompt Studio IDE 与 ETL 工作流中的 LLM 提取、向量索引与 Prompt 执行任务并通过回调与 Socket.IO 将结果实时推回浏览器。阅读后你将掌握 Executor Worker 的消息链路、队列与并发配置、Docker 部署方式、本地开发启动方法以及execute_extraction任务与LegacyExecutor操作分发机制的源码级工作原理。一、Executor Worker 是什么Executor Worker 是 Unstract 平台中负责重活的 Celery 工作进程它监听消息队列接收包含ExecutionContext字典的任务通过执行器框架Executor Framework运行 LLM 提取、文档索引、Prompt 回答等操作最终返回一个ExecutionResult字典。官方 README 对其定位只有一句话Celery worker that handles LLM extraction, indexing, and prompt execution——即处理 LLM 提取、索引和 Prompt 执行的 Celery 工作进程。它在整个 Unstract 平台架构中处于执行引擎的位置所有需要调用大模型LLM能力完成非结构化数据抽取的操作最终都会汇聚到这个工作进程上来。二、端到端工作原理一次Run的完整旅程2.1 总体消息链路官方 README 用一条 ASCII 链路概括了 Executor Worker 在系统中的位置Browser → Django Backend → RabbitMQ → Executor Worker → Callback → WebSocket → Browser2.2 五个步骤拆解README 将一次完整的 Prompt Studio 运行拆解为 5 个步骤用户点击 Run用户在 Prompt Studio IDE 中点击运行Django 后端将任务分发到celery_executor_legacy队列Executor worker 消费任务执行器工作进程从队列中拾取任务运行 LLM 提取触发回调执行结果触发ide_callback队列上的回调IDE 回调 worker 持久化结果ide_callback工作进程通过内部 API 持久化结果并通过 Socket.IO 推送浏览器实时接收浏览器实时收到结果。2.3 涉及的服务清单README 给出了参与该流程的五个服务及其职责ServicePurposeworker-executor-v2Runs LLM extraction, indexing, promptsworker-ide-callbackPost-execution callbacks via internal API Socket.IO eventsbackendDjango REST API Socket.IOplatform-serviceAdapter credential managementprompt-servicePrompt template service从仓库当前的 docker-compose.yaml 可以看到worker-ide-callback是 Celery 形态的回调消费者而新增的worker-pg-ide-callback位于 docker-compose 约 882 行起则是 PG 队列形态的回调消费者二者职责相同回调只做内部 API 写入和 WebSocket 事件推送不再向 Celery 发起后续分发因此非常轻量默认并发数为 2PG_IDE_CALLBACK_CONCURRENCY可见度超时visibility timeout默认 120 秒健康检查过期窗口默认 180 秒。三、任务入口execute_extraction 与 ExecutionContext3.1 唯一的 Celery 任务入口workers/executor/tasks.py 中定义了 Executor Worker 的核心任务execute_extraction其 docstring 明确指出这是所有提取操作的唯一 Celery 任务入口工作流路径structure tool task和 IDE 路径PromptStudioHelper都会分发到这个任务。任务签名与重试策略如下worker_task( bindTrue, nameTaskName.EXECUTE_EXTRACTION, autoretry_for(ConnectionError, TimeoutError, OSError), retry_backoffTrue, retry_backoff_max60, max_retries3, retry_jitterTrue, ) def execute_extraction(self, execution_context_dict: dict) - dict:入参序列化后的ExecutionContext字典出参序列化后的ExecutionResult字典针对ConnectionError、TimeoutError、OSError三类瞬时错误自动重试最多 3 次退避上限 60 秒并带抖动jitter避免重试风暴。3.2 任务执行的核心流程execute_extraction的执行流程可以归纳为四步反序列化上下文ExecutionContext.from_dict(execution_context_dict)解析执行上下文若上下文非法KeyError/ValueError直接返回ExecutionResult.failure(...)日志关联根据operation如ide_index、structure_pipeline从嵌套参数中提取tool_id、run_id、doc_name等字段构建日志组件_log_component用于把前端流式日志与实际执行对应起来编排执行创建ExecutionOrchestrator调用orchestrator.execute(context)得到ExecutionResult用量上报若结果元数据中带usage_records通过UsageAPIClient.bulk_create_usage(...)批量上报计费用量上报失败会以 ERROR 级别记录日志保证计费数据可从日志恢复。源码中还定义了一个有趣的白名单集合_LLM_BEARING_OPS frozenset( { answer_prompt, single_pass_extraction, summarize, structure_pipeline, } )当这些必然产生 LLM 调用的操作成功却没有产出任何 usage 记录时任务会记录一条提示性日志用于尽早暴露计费数据缺失的问题。3.3 执行上下文与结果ExecutionContext与ExecutionResult均来自unstract.sdk1.execution包见 unstract/sdk1/src 下的execution/context.py、execution/result.py。从任务与执行器代码可以推断ExecutionContext至少携带以下关键字段executor_name选用哪个执行器operation执行哪种操作见下文 Operation 枚举run_id/request_id链路追踪标识execution_source执行来源ide或toolorganization_id、execution_id、file_execution_id组织与文件执行归属log_events_id日志流事件 IDexecutor_params具体操作的参数字典。3.4 健康检查任务workers/executor/worker.py 还注册了名为executor_health的自定义健康检查它从ExecutorRegistry.list_executors()列出已注册的执行器返回健康状态、已注册执行器数量与监听的队列列表监控系统可以直接消费该健康检查结果判断执行器是否就绪。若注册表查询失败健康状态降级为DEGRADED并携带错误信息。四、队列与并发配置4.1 监听队列README 明确Executor Worker 监听celery_executor_legacy队列可通过CELERY_QUEUES_EXECUTOR环境变量配置。在 docker-compose.yaml 的worker-executor-v2服务中实际默认值已经扩展为三个队列- CELERY_QUEUES_EXECUTOR${CELERY_QUEUES_EXECUTOR:-celery_executor_legacy,celery_executor_agentic,celery_executor_agentic_table}即默认同时监听celery_executor_legacy、celery_executor_agentic、celery_executor_agentic_table三个队列其中后两个服务于 Agentic智能体式提取场景。同一队列名也出现在 workers/run-worker.sh 的 PG 队列角色pg-executor中executor;celery_executor_legacy,celery_executor_agentic,celery_executor_agentic_table说明 PG 队列形态的执行器消费者与 Celery 形态的执行器保持队列一致。4.2 关键环境变量README 给出的配置变量表是 Executor Worker 调优的核心依据VariableDefaultDescriptionWORKER_EXECUTOR_CONCURRENCY2Number of concurrent executor processesWORKER_EXECUTOR_POOLpreforkCelery pool typeEXECUTOR_TASK_TIME_LIMIT3600Hard timeout per task (seconds)EXECUTOR_TASK_SOFT_TIME_LIMIT3300Soft timeout per task (seconds)EXECUTOR_RESULT_TIMEOUT3600How long callers wait for resultsEXECUTOR_AUTOSCALE2,1Max,min worker autoscale这些变量的底层落地位置如下均可在仓库中逐一验证WORKER_EXECUTOR_CONCURRENCY默认 2在 docker-compose.yaml 中映射为CELERY_CONCURRENCY${WORKER_EXECUTOR_CONCURRENCY:-2}即 Celery--concurrency参数WORKER_EXECUTOR_POOL默认prefork映射为CELERY_POOL${WORKER_EXECUTOR_POOL:-prefork}即 Celery--pool参数WORKER_EXECUTOR_PREFETCH_MULTIPLIER默认 1映射为CELERY_PREFETCH_MULTIPLIER控制每个进程预取消息数WORKER_EXECUTOR_EXTRA_ARGS透传给 Celery 的额外参数CELERY_EXTRA_ARGSEXECUTOR_TASK_TIME_LIMIT/EXECUTOR_TASK_SOFT_TIME_LIMIT默认 3600/3300 秒在 workers/sample.env270–275 行中随EXECUTOR_WORKER_NAMEexecutor-worker、EXECUTOR_HEALTH_PORT8088、EXECUTOR_AUTOSCALE2,1、EXECUTOR_RESULT_TIMEOUT3600一起提供软超时3300s略小于硬超时3600s便于在硬杀前留出收尾窗口EXECUTOR_AUTOSCALE2,1max,min 自动伸缩配置2 个最大进程、1 个最小进程。说明LLM 推理耗时长执行器任务可能运行接近一小时docker-compose 中 PG 队列形态的执行器消费者worker-pg-executor将可见度超时设置为WORKER_PG_EXECUTOR_VT_SECONDS:-3660明确高于执行器单任务硬限制EXECUTOR_TASK_TIME_LIMIT3600s防止兄弟副本在任务运行中重复认领导致重复执行。五、Docker 部署README 指出Executor Worker 在 docker/docker-compose.yaml 中定义为worker-executor-v2服务使用统一 worker 镜像unstract/worker-unified并以executor命令启动。对应的 compose 片段约 511–538 行要点如下worker-executor-v2: image: unstract/worker-unified:${VERSION} container_name: unstract-worker-executor-v2 restart: unless-stopped command: [executor] ports: - 8092:8088 env_file: - ../workers/.env - ./essentials.env depends_on: - db - redis - rabbitmq - platform-service environment: - APPLICATION_NAMEunstract-worker-executor-v2 - WORKER_TYPEexecutor - WORKER_NAMEexecutor-worker-v2 - EXECUTOR_METRICS_PORT8088 - HEALTH_PORT8088 - CELERY_QUEUES_EXECUTOR${CELERY_QUEUES_EXECUTOR:-celery_executor_legacy,celery_executor_agentic,celery_executor_agentic_table} - CELERY_POOL${WORKER_EXECUTOR_POOL:-prefork} - CELERY_PREFETCH_MULTIPLIER${WORKER_EXECUTOR_PREFETCH_MULTIPLIER:-1} - CELERY_CONCURRENCY${WORKER_EXECUTOR_CONCURRENCY:-2} - CELERY_EXTRA_ARGS${WORKER_EXECUTOR_EXTRA_ARGS:-}要点解读镜像与命令统一镜像unstract/worker-unified容器启动命令为executor即启动执行器入口依赖关系依赖db、redis、rabbitmq、platform-service数据库、缓存/消息、适配器凭据管理健康与指标端口容器内 8088 端口同时作为健康检查与指标端口宿主映射为 8092卷挂载挂载./workflow_data:/data与TOOL_REGISTRY_CONFIG_SRC_PATH指向的工具注册表配置以及prompt_studio_data数据卷无额外配置要求README 特别说明——执行器随./run-platform.sh自动启动开箱即用无需额外配置The executor worker starts automatically with ./run-platform.sh — no extra configuration needed。此外worker-unified.Dockerfile 显示统一镜像构建时会扫描并安装unstract.executor.executors执行器类如 table_extractor与unstract.executor.plugins工具插件如 highlight-data、challenge等执行器插件这解释了为何执行器容器天然具备可插拔能力。六、本地开发快速启动 Executor WorkerREADME 给出本地开发的最小启动流程cd workers cp sample.env .env # Edit .env: change Docker hostnames to localhost ./run-worker.sh executor结合 workers/run-worker.sh 的完整用法实际开发中还可以使用更多选项# 后台运行daemon 模式日志写入 workers/executor/executor.log ./run-worker.sh -d executor # 指定并发数、日志级别与队列 ./run-worker.sh -c 4 -l DEBUG executor # 指定 Celery pool 类型prefork/threads/gevent/solo/eventlet ./run-worker.sh -P prefork executor # 查看运行状态 / 实时日志 / 重启 ./run-worker.sh -s ./run-worker.sh -L executor ./run-worker.sh -r executor该脚本对executor的默认处理包括工作目录为workers/executor启动命令形如uv run celery -A worker worker --loglevelinfo --queuescelery_executor_legacy --hostnameexecutor-worker%h --concurrency2默认并发 2与WORKER_EXECUTOR_CONCURRENCY一致默认健康端口 8088EXECUTOR_HEALTH_PORT环境校验要求至少具备INTERNAL_SERVICE_API_KEY、INTERNAL_API_BASE_URL、CELERY_BROKER_BASE_URL、DB_HOST、DB_USER、DB_PASSWORD、DB_NAME等变量。七、执行器注册表与 LegacyExecutor操作分发机制7.1 入口处的导入策略workers/executor/worker.py 末尾有两条关键导入import executor.executors # noqa: E402, F401 import executor.tasks # noqa: E402, F401注释解释了原因导入executor.executors是为了触发ExecutorRegistry.register装饰器在导入期执行从而把内置执行器注册进注册表。同时 workers/executor/tasks.py 顶部也导入了executor.executors因为 PG 队列消费端通过根 worker 导入加载的是executor/tasks.py而非executor/worker.py——若此处不导入消费端会报 No executor registered。Python 模块缓存保证了重复导入无害幂等。7.2 内置执行器LegacyExecutorworkers/executor/executors/legacy_executor.py 中的LegacyExecutor是仓库内置的核心执行器以ExecutorRegistry.register装饰器注册name属性为legacy。其 docstring 说明它包装了完整的 prompt-service 提取流水线将ExecutionContext请求按Operation枚举路由到对应的处理方法每个处理方法对应原 prompt-service 的一个 HTTP 端点。_OPERATION_MAP: dict[str, str] { Operation.EXTRACT.value: _handle_extract, Operation.INDEX.value: _handle_index, Operation.ANSWER_PROMPT.value: _handle_answer_prompt, Operation.SINGLE_PASS_EXTRACTION.value: _handle_single_pass_extraction, Operation.SUMMARIZE.value: _handle_summarize, Operation.IDE_INDEX.value: _handle_ide_index, Operation.STRUCTURE_PIPELINE.value: _handle_structure_pipeline, }7.3 七种操作一览基于操作映射与各处理方法的实现可以整理出执行器支持的操作清单Operation处理方法职责extract_handle_extract通过 x2text 适配器做文本提取index_handle_index向量库索引embedding vector DBanswer_prompt_handle_answer_prompt多 Prompt 提取变量替换、检索、LLM 补全、类型化后处理single_pass_extraction_handle_single_pass_extraction单次读取文件、构建合并 Prompt、单次 LLM 调用summarize_handle_summarize文档摘要ide_index_handle_ide_index复合操作提取 索引一步完成structure_pipeline_handle_structure_pipeline复合操作提取 → 摘要 → 索引 → 回答完整流水线execute()方法的错误处理契约值得注意任何LegacyExecutorError子类异常都会被捕获并映射为ExecutionResult.failure()保证调用方永远得到一个结果对象而非裸异常失败时若存在log_events_id还会通过 shim 把错误实时流式推送到前端。同时中途失败产生的部分 usage 记录partial_usage_records会被合并进失败结果的元数据避免计费数据丢失。7.4 各操作的执行要点源码级extract文本提取要求x2text_instance_id与file_path两个必填参数可选output_file_path、enable_highlight、tags、usage_kwargs、tool_execution_metadata、execution_data_dir等通过X2Text构建 x2text 适配器实例并执行process()结果数据包含extracted_text若适配器返回行级元数据还会附带highlight_metadata供智能体提取做 PDF 源引用。index向量索引必填embedding_instance_id、vector_db_instance_id、x2text_instance_id、file_path默认chunk_size512、chunk_overlap128当chunk_size 0时跳过向量化仅生成doc_idIndexingUtils.generate_index_key当文档已索引且reindexFalse时跳过重复索引。索引的 embedding 用量在成功与失败路径都会被 flush 到 usage 记录。answer_prompt多 Prompt 提取payload 包含tool_settings、outputs每个输出含name/prompt/chunk-size等与file_path等核心步骤为变量替换VariableReplacementService→ 按检索策略取上下文RetrievalService支持 simple/subquestion/fusion/recursive/router/keyword_table/automerging→ LLM 补全 → 类型化转换NUMBER/EMAIL/DATE/BOOLEAN/JSON/TEXT→ 可选 lookup 富化、Webhook 后处理、challenge 校验、evaluation 评估。tool_settings中的 chunk/LLM/向量库等字段会作为默认值合并进每个 output保证单遍single-pass与多遍answer_prompt两种 payload 形态都能工作。single_pass_extraction优先委托给云插件ExecutorRegistry.get(single_pass_extraction)读取文件一次、构建一个合并 Prompt、发起一次 LLM 调用若插件未安装则回退到_handle_answer_prompt。summarize摘要入参为llm_adapter_instance_id、summarize_prompt、context全文、prompt_keys重点字段构造提示词 重点字段 上下文的摘要 Prompt 后调用 LLM 补全。在流水线中摘要结果会被写入文件系统缓存命中缓存时不再触发 LLM 调用。ide_index复合提取 索引一次执行器调用内完成_handle_extract与_handle_index消除了后端 Celery worker 在两个步骤间的阻塞等待若index_params已带extracted_textmarker 命中则跳过提取步骤支持可选的 summarize 步骤子步骤的 usage 记录会被吸收合并到最终结果。structure_pipeline复合完整流水线在一次执行器调用内依次执行 提取 → 摘要可选→ 逐输出索引带参数去重→ 回答/单遍提取替代了原本三次串行的dispatcher.dispatch()调用避免长时间占用 file_processing worker 槽位pipeline_options中的skip_extraction_and_indexing、is_summarization_enabled、is_single_pass_enabled等开关控制流水线分支最终合并output、metadata、metrics返回结构化结果。7.5 可插拔执行器与工具插件从 workers/executor/executors/init.py 可以看到执行器包在导入时会注册内置的LegacyExecutor通过ExecutorPluginLoader.discover_executors()按 entry points 发现并注册云端执行器插件未安装时返回空列表。运行期还支持按需获取插件执行器例如_run_table_extraction使用ExecutorRegistry.get(table)处理 TABLE/RECORD 类型输出_run_line_item_extraction使用ExecutorRegistry.get(line_item)处理 LINE_ITEM 类型输出——若对应插件未安装会明确报错提示Install the table_extractor plugin或Install the line_item_extractor plugin。highlight-data、challenge、evaluation 等工具插件同理均通过 workers/executor/executors/plugins 下的加载器动态获取。八、安全与隐私细节workers/executor/worker.py 中包含一个值得关注的安全实践# Suppress Celery trace logging of task return values. # The trace logger prints the full result dict on task success, which # can contain sensitive customer data (extracted text, summaries, etc.). logging.getLogger(celery.app.trace).setLevel(logging.WARNING)Celery 的 trace 日志会在任务成功时打印完整的返回值字典而执行器任务的结果可能包含提取文本、摘要等客户敏感数据因此执行器在启动时显式将celery.app.trace日志级别提升到 WARNING避免敏感结果被写入日志。九、常见问题与排障思路任务一直不执行确认执行器容器/进程在运行且CELERY_QUEUES_EXECUTOR包含任务实际投递的队列默认三个celery_executor_legacy,celery_executor_agentic,celery_executor_agentic_tableNo executor registered任务模块未触发执行器注册检查是否导入了executor.executors该导入幂等重复无害提取超时检查EXECUTOR_TASK_TIME_LIMIT硬超时 3600s与EXECUTOR_TASK_SOFT_TIME_LIMIT软超时 3300s是否与你的长文档处理场景匹配并同步调整 PG 消费端的可见度超时需高于硬超时如 3660s计费/用量数据缺失观察执行器日志中关于usage_records的 ERROR 日志bulk_create_usage returned failure ...与 LLM-bearing 操作无用量记录的告警日志TABLE/LINE_ITEM 报插件缺失需要安装对应的执行器插件table_extractor / line_item_extractor统一镜像构建时会自动扫描安装unstract.executor.executors与unstract.executor.plugins目录下的插件。十、延伸阅读执行器入口与健康检查workers/executor/worker.py任务定义与用量上报workers/executor/tasks.py核心执行器实现workers/executor/executors/legacy_executor.py执行器包与插件发现workers/executor/executors/init.py、workers/executor/executors/plugins执行上下文与结果定义unstract/sdk1/src下的 execution 模块队列与并发环境变量workers/sample.env、docker/docker-compose.yaml本地启动脚本workers/run-worker.sh统一 worker 镜像构建docker/dockerfiles/worker-unified.Dockerfile回调侧消费端worker-ide-callback/worker-pg-ide-callback见 docker/docker-compose.yaml【免费下载链接】unstractLLM-Driven Extraction of Unstructured Data — Built for API Deployments ETL Pipeline Workflows项目地址: https://gitcode.com/GitHub_Trending/un/unstract创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表