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 的实现依赖requests与xmltodict两个第三方库(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(调度用)、pandas、pillow(表格与图像任务用)等,如果你的工作流只用 ServiceTask,也可以单独安装requests与xmltodict两个库。
快速上手: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 中剩余字段(url、method、params、batch、extract等)作为关键字参数传给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):各参数语义如下表:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
url | str | None | 要连接的远程服务地址 |
method | str | None | HTTP 方法,"get"或"post";不传或非 get 时按 POST 处理 |
params | dict | None | 默认查询/请求参数;值为None的键会用当前数据元素动态填充 |
batch | bool | True | 为True时所有元素在一次批量请求中发送;为False时每个元素单独发起一次请求 |
extract | str / list | None | 需要从响应中提取的字段,支持单个字符串或层级列表 |
其中extract在register中被标准化为列表(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(正则过滤输入)、unpack、merge、concurrency、onetomany等,具体可查阅 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 处理与输出合并流程(如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, 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-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"]。测试用例中服务返回<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。
三种场景分别验证了:
- GET + 动态 params:
params.text为None,被替换为工作流输入元素,请求走 query string; - POST + 无 params:
params为空,直接以工作流数据作为 JSON 请求体发送; - 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),因此下游任务拿到的是按原输入顺序对齐的结果,可以无缝衔接。
注意事项小结
- 必须安装
requests与xmltodict,否则构造 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),仅供参考