txtai ServiceTask 详解:在 Workflow 中调用 HTTP 服务并解析响应
2026/9/15 13:05:31 网站建设 项目流程

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/txtai

ServiceTask 是 txtai 工作流框架中用于对接外部 HTTP 服务的任务类型:它把工作流中的每个数据元素作为请求参数发送到远程 URL,并自动把 JSON / XML 响应解析回结构化数据,供下游任务继续处理。本文基于仓库中 service.md 文档,结合 service.py 源码与 testapiworkflow.py 测试用例,完整讲解 ServiceTask 的安装前提、Python / YAML 两种创建方式、全部构造参数语义以及底层请求执行细节,帮助你快速把任意 HTTP 服务接入 txtai 工作流。

什么是 ServiceTask

在 txtai 的 Workflow 体系中,工作流由一串可调用的 Task 组成,数据以流式批次的方式依次经过每个任务(参见 工作流总览 与 Task 基类)。ServiceTask 是其中的一个特殊任务,它的职责很单一:向远程 HTTP 服务发起请求,并把响应内容作为该任务的处理结果返回

从源码类定义看(service.py):

class 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 的实现依赖requestsxmltodict两个第三方库(service.py):

# Conditional import try: import requests import xmltodict XML_TO_DICT = True except ImportError: XML_TO_DICT = False

如果这两个库未安装,实例化 ServiceTask 会直接抛出ImportError(service.py):

raise ImportError('ServiceTask is not available - install "workflow" extra to enable')

这两个库正是由 txtai 的workflow扩展包提供(setup.py):

extras["workflow"] = [ "apache-libcloud>=3.3.1", "croniter>=1.2.0", "openpyxl>=3.0.9", "pandas>=1.1.0", "pillow>=7.1.2", "requests>=2.26.0", "xmltodict>=0.12.0", ]

因此使用前请先安装:

pip install txtai[workflow]

说明:workflowextra 还包含croniter(调度用)、pandaspillow(表格与图像任务用)等,如果你的工作流只用 ServiceTask,也可以单独安装requestsxmltodict两个库。

快速上手:Python 方式

原文档给出了最简用法(service.md):

from txtai.workflow import ServiceTask, Workflow workflow = Workflow([ServiceTask(url="https://service.url/action)]) workflow(["parameter"])

这段代码的含义是:工作流接收["parameter"]这样的输入列表,ServiceTask 把"parameter"作为请求参数发送到https://service.url/action,服务端返回的 JSON / XML 被解析后作为工作流输出。

一个更完整的写法如下:

from txtai.workflow import ServiceTask, Workflow workflow = Workflow([ ServiceTask( url="https://api.example.com/search", method="get", params={"q": None, "limit": 10}, # 值为 None 的参数会被工作流数据填充 batch=True, # 所有元素一次请求 ) ]) # 注意:Workflow 以生成器方式返回结果,需要用 list() 或 for 消费 for result in workflow(["txtai", "semantic search"]): print(result)

关于Workflow返回生成器、需要显式消费输出的行为,以及批量流式处理机制,可参考 工作流总览。

配置驱动:YAML 方式

ServiceTask 同样可以通过工作流配置创建(service.md):

workflow: 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 中剩余字段(urlmethodparamsbatchextract等)作为关键字参数传给ServiceTask(**config)(factory.py),最终落到register方法上完成初始化。

在 API 配置中,这种组合常与schedule搭配,实现周期性拉取外部服务数据后接续索引(configuration.md):

workflow: index: schedule: cron: 0/10 * * * * * elements: ["api params"] tasks: - task: service url: api url - action: index

构造参数详解(register 源码级解读)

ServiceTask 的全部自定义参数定义在register方法中(service.py),方法签名如下:

def register(self, url=None, method=None, params=None, batch=True, extract=None):

各参数语义如下表:

参数类型默认值说明
urlstrNone要连接的远程服务地址
methodstrNoneHTTP 方法,"get""post";不传或非 get 时按 POST 处理
paramsdictNone默认查询/请求参数;值为None的键会用当前数据元素动态填充
batchboolTrueTrue时所有元素在一次批量请求中发送;为False时每个元素单独发起一次请求
extractstr / listNone需要从响应中提取的字段,支持单个字符串或层级列表

其中extractregister中被标准化为列表(service.py):

self.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(正则过滤输入)、unpackmergeconcurrencyonetomany等,具体可查阅 base.py 与 task/index.md。

批量与逐元素执行

execute方法根据batch决定请求的组织方式(service.py):

def execute(self, elements, executor=None): if self.batch: elements = self.request(elements) else: elements = [self.request(element) for element in elements] return super().execute(elements, executor)
  • batch=True(默认):整批元素一次性传给request,适合服务端支持批量入参的接口,请求次数少、效率高;
  • batch=False:每个元素单独调用一次服务,适合逐条处理、或需要单独跟踪每次调用结果的场景。

执行完请求后,结果会交给基类的execute继续走后续的 action 处理与输出合并流程(如mergepostprocess等,见 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, params=params) else: response = requests.post(self.url, json=params)
  • method="get"(不区分大小写):走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-Typexml结尾(如application/xmltext/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"]。测试用例中服务返回<row><text>test</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.py):

  • GET /testget返回 JSON 数组[{"text": "test"}]Content-Type: application/json
  • POST /testpost读取 JSON 请求体,按.切分文本后返回嵌套 JSON 数组;
  • GET /xml返回 XML 文档<row><text>test</text></row>Content-Type: application/xml

三种场景分别验证了:

  1. GET + 动态 paramsparams.textNone,被替换为工作流输入元素,请求走 query string;
  2. POST + 无 paramsparams为空,直接以工作流数据作为 JSON 请求体发送;
  3. GET + batch=False + extract:关闭批量模式逐元素请求,XML 响应经xmltodict解析后,再按extract: row提取出{"text": "test"}字典。

对应测试方法testServiceGet(testapiworkflow.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),因此下游任务拿到的是按原输入顺序对齐的结果,可以无缝衔接。

注意事项小结

  • 必须安装requestsxmltodict,否则构造 ServiceTask 即抛ImportError,提示安装workflowextra;
  • method缺省时走 POST 且请求体为 JSON,需要 query 传参时务必显式设置method: get
  • params中值为空的键才会被动态数据填充,其余键作为常量发送;
  • 响应解析依赖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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询