txtai ServiceTask 详解:在 Workflow 中调用 HTTP 服务并解析响应
txtai ServiceTask 详解在 Workflow 中调用 HTTP 服务并解析响应【免费下载链接】txtai All-in-one AI framework for semantic search, LLM orchestration and language model workflows项目地址: https://gitcode.com/GitHub_Trending/tx/txtaiServiceTask 是 txtai 工作流框架中用于对接外部 HTTP 服务的任务类型它把工作流中的每个数据元素作为请求参数发送到远程 URL并自动把 JSON / XML 响应解析回结构化数据供下游任务继续处理。本文基于仓库中 service.md 文档结合 service.py 源码与 testapiworkflow.py 测试用例完整讲解 ServiceTask 的安装前提、Python / YAML 两种创建方式、全部构造参数语义以及底层请求执行细节帮助你快速把任意 HTTP 服务接入 txtai 工作流。什么是 ServiceTask在 txtai 的 Workflow 体系中工作流由一串可调用的 Task 组成数据以流式批次的方式依次经过每个任务参见 工作流总览 与 Task 基类。ServiceTask 是其中的一个特殊任务它的职责很单一向远程 HTTP 服务发起请求并把响应内容作为该任务的处理结果返回。从源码类定义看service.pyclass ServiceTask(Task): Task to runs requests against remote service urls. 它继承自Task基类因此天然具备工作流任务的全部通用能力——select过滤、unpack解包、merge合并、concurrency并发等详见 base.py。ServiceTask 自己额外承担的是通过register方法注册一组与 HTTP 请求相关的参数。典型应用场景包括把 txtai 工作流的中间结果发给外部 API如第三方 NLP 服务、内部推理服务并取回结果在 定时调度的工作流 中周期性拉取接口数据再接续到索引、摘要等下游任务对接返回 JSON 或 XML 的遗留系统由 ServiceTask 完成响应解析与字段提取。安装前提workflow 扩展依赖ServiceTask 的实现依赖requests与xmltodict两个第三方库service.py# Conditional import try: import requests import xmltodict XML_TO_DICT True except ImportError: XML_TO_DICT False如果这两个库未安装实例化 ServiceTask 会直接抛出ImportErrorservice.pyraise ImportError(ServiceTask is not available - install workflow extra to enable)这两个库正是由 txtai 的workflow扩展包提供setup.pyextras[workflow] [ apache-libcloud3.3.1, croniter1.2.0, openpyxl3.0.9, pandas1.1.0, pillow7.1.2, requests2.26.0, xmltodict0.12.0, ]因此使用前请先安装pip install txtai[workflow]说明workflowextra 还包含croniter调度用、pandas、pillow表格与图像任务用等如果你的工作流只用 ServiceTask也可以单独安装requests与xmltodict两个库。快速上手Python 方式原文档给出了最简用法service.mdfrom txtai.workflow import ServiceTask, Workflow workflow Workflow([ServiceTask(urlhttps://service.url/action)]) workflow([parameter])这段代码的含义是工作流接收[parameter]这样的输入列表ServiceTask 把parameter作为请求参数发送到https://service.url/action服务端返回的 JSON / XML 被解析后作为工作流输出。一个更完整的写法如下from txtai.workflow import ServiceTask, Workflow workflow Workflow([ ServiceTask( urlhttps://api.example.com/search, methodget, params{q: None, limit: 10}, # 值为 None 的参数会被工作流数据填充 batchTrue, # 所有元素一次请求 ) ]) # 注意Workflow 以生成器方式返回结果需要用 list() 或 for 消费 for result in workflow([txtai, semantic search]): print(result)关于Workflow返回生成器、需要显式消费输出的行为以及批量流式处理机制可参考 工作流总览。配置驱动YAML 方式ServiceTask 同样可以通过工作流配置创建service.mdworkflow: name: tasks: - task: service url: https://service.url/action其中task: service是任务类型标识。txtai 的任务工厂会把该字符串解析为ServiceTask从源码看工厂在遇到不带包名的任务名时会自动拼接为txtai.workflow.task.TaskName并加载对应类factory.py# Local task if no package if . not in task: # Get parent package task ..join(__name__.split(.)[:-1]) . task.capitalize() Task随后把 YAML 中剩余字段url、method、params、batch、extract等作为关键字参数传给ServiceTask(**config)factory.py最终落到register方法上完成初始化。在 API 配置中这种组合常与schedule搭配实现周期性拉取外部服务数据后接续索引configuration.mdworkflow: index: schedule: cron: 0/10 * * * * * elements: [api params] tasks: - task: service url: api url - action: index构造参数详解register 源码级解读ServiceTask 的全部自定义参数定义在register方法中service.py方法签名如下def register(self, urlNone, methodNone, paramsNone, batchTrue, extractNone):各参数语义如下表参数类型默认值说明urlstrNone要连接的远程服务地址methodstrNoneHTTP 方法get或post不传或非 get 时按 POST 处理paramsdictNone默认查询/请求参数值为None的键会用当前数据元素动态填充batchboolTrue为True时所有元素在一次批量请求中发送为False时每个元素单独发起一次请求extractstr / listNone需要从响应中提取的字段支持单个字符串或层级列表其中extract在register中被标准化为列表service.pyself.extract extract if self.extract: self.extract [self.extract] if isinstance(self.extract, str) else self.extract也就是说传row与传[row]等价传[result, items]则表示按层级逐层向下取字段见下文 request 部分。补充作为Task子类ServiceTask 同样接受基类的通用参数如select正则过滤输入、unpack、merge、concurrency、onetomany等具体可查阅 base.py 与 task/index.md。批量与逐元素执行execute方法根据batch决定请求的组织方式service.pydef execute(self, elements, executorNone): if self.batch: elements self.request(elements) else: elements [self.request(element) for element in elements] return super().execute(elements, executor)batchTrue默认整批元素一次性传给request适合服务端支持批量入参的接口请求次数少、效率高batchFalse每个元素单独调用一次服务适合逐条处理、或需要单独跟踪每次调用结果的场景。执行完请求后结果会交给基类的execute继续走后续的 action 处理与输出合并流程如merge、postprocess等见 base.py。请求执行细节request 源码级解读request是 ServiceTask 的核心方法完整实现了参数合并、方法选择、响应解析与字段提取service.py。逐段拆解如下。1. 动态参数合并if not self.params: params data else: # Create copy of parameters params self.params.copy() # Add data to parameters for key in params: if not params[key]: params[key] data若未配置params则直接把工作流数据data作为请求参数若配置了params先拷贝一份再把其中所有值为空的键None、空字符串、0、空列表等均视为 falsy替换为当前数据data。例如配置params: {q: None, limit: 10}则limit保持常量10而q被填充为当前元素。测试用例中的params: {text: }值为None正是利用了这一机制testapiworkflow.py。2. HTTP 方法与请求发送# Run request if self.method and self.method.lower() get: response requests.get(self.url, paramsparams) else: response requests.post(self.url, jsonparams)methodget不区分大小写走requests.get参数以 URL query string 形式携带其他情况包括不传method一律走requests.post参数以 JSON 请求体发送。3. 按 Content-Type 解析响应# Parse data based on content-type mimetype response.headers[Content-Type].split(;)[0] if mimetype.lower().endswith(xml): data xmltodict.parse(response.text) else: data response.json()服务端响应头Content-Type以xml结尾如application/xml、text/xml时使用xmltodict.parse把 XML 文本解析为字典其余情况如application/json使用response.json()解析为 Python 对象。这解释了为什么 ServiceTask 需要xmltodict依赖——它负责 XML 响应到字典的转换。4. 响应字段提取# Extract content from response, if necessary if self.extract: for tag in self.extract: data data[tag]配置了extract时会按顺序逐层取字段例如extract[result, items]等价于data data[result][items]。测试用例中服务返回rowtexttest/text/row配合extract: row即可直接取到{text: test}见下文。测试用例验证仓库中的 API 工作流测试testapiworkflow.py覆盖了 ServiceTask 的三种典型形态是理解各参数组合行为的最佳参考get: tasks: - task: service url: http://127.0.0.1:8001/testget method: get params: text: post: tasks: - task: service url: http://127.0.0.1:8001/testpost params: xml: tasks: - task: service url: http://127.0.0.1:8001/xml method: get batch: false extract: row params: text:对应测试服务的行为testapiworkflow.pyGET /testget返回 JSON 数组[{text: test}]Content-Type: application/jsonPOST /testpost读取 JSON 请求体按.切分文本后返回嵌套 JSON 数组GET /xml返回 XML 文档rowtexttest/text/rowContent-Type: application/xml。三种场景分别验证了GET 动态 paramsparams.text为None被替换为工作流输入元素请求走 query stringPOST 无 paramsparams为空直接以工作流数据作为 JSON 请求体发送GET batchFalse extract关闭批量模式逐元素请求XML 响应经xmltodict解析后再按extract: row提取出{text: test}字典。对应测试方法testServiceGettestapiworkflow.py通过 API 触发get工作流并断言结果长度直接印证了上述执行链路。与下游任务编排ServiceTask 的结果会作为下一个任务的输入继续流转。比如把外部服务的响应接续到标签分类、摘要或嵌入索引workflow: enrich: tasks: - task: service url: https://service.url/action method: get params: text: - action: labels args: [[positive, negative]]也可以像 configuration.md 中的示例那样在服务任务之后接action: index直接入库workflow: index: schedule: cron: 0/10 * * * * * elements: [api params] tasks: - task: service url: api url - action: index从源码结构看ServiceTask 返回的解析后数据字典、列表或标量会进入基类execute的输出处理管线再以filteredpack/pack的方式与原始元素重新组装base.py因此下游任务拿到的是按原输入顺序对齐的结果可以无缝衔接。注意事项小结必须安装requests与xmltodict否则构造 ServiceTask 即抛ImportError提示安装workflowextramethod缺省时走 POST 且请求体为 JSON需要 query 传参时务必显式设置method: getparams中值为空的键才会被动态数据填充其余键作为常量发送响应解析依赖Content-Type头以xml结尾按 XML 解析否则按 JSON 解析请确保服务端正确返回该响应头extract支持层级提取多层结构传列表即可按序下钻工作流输出是生成器使用时需list()或for循环消费请求才会真正执行。【免费下载链接】txtai All-in-one AI framework for semantic search, LLM orchestration and language model workflows项目地址: https://gitcode.com/GitHub_Trending/tx/txtai创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考