txtai 工作流 Task 完全指南:可调用处理单元、多动作并发与列合并
2026/9/15 14:11:49 网站建设 项目流程

txtai 工作流 Task 完全指南:可调用处理单元、多动作并发与列合并

【免费下载链接】txtai💡 All-in-one AI framework for semantic search, LLM orchestration and language model workflows项目地址: https://gitcode.com/GitHub_Trending/tx/txtai

工作流(Workflow)是 txtai 中承载大规模数据流处理的核心构造,而 Task 则是构成工作流的原子处理单元——它接收可迭代的数据元素,对每个元素执行一个或多个动作(action),再输出转换后的结果。本文以 docs/workflow/task/index.md 为骨架,结合 Task 源码 与 测试用例,系统讲解 Task 的参数体系、与 Pipeline/Workflow 的协作方式、多动作并发(multithreading/multiprocessing)以及 hstack / vstack / concat 三种合并模式,读完即可在 Python 与 YAML 配置两种方式下熟练构建自己的任务链。

Task 是什么

在 txtai 中,Workflows execute tasks(工作流执行任务)。Task 是带有若干控制参数的可调用对象(callable objects),负责在流程的某个步骤控制数据的处理方式。它与 Pipeline 的定位有明显区别:

While similar to pipelines, tasks encapsulate processing and don't perform significant transformations on their own. Tasks perform logic to prepare content for the underlying action(s).

也就是说,Task 本身不做重量级变换,它更像一个"编排壳":负责准备数据、过滤数据、选择列、调度并发、合并多路输出,而真正的变换逻辑由底层的 action(可调用对象,常常就是 Pipeline 或普通函数)完成。这与 Workflow 文档 中"Workflows are a simple yet powerful construct that takes a callable and returns elements"的描述一脉相承。

从源码看,Task 基类 是"所有工作流任务的基类"(Base class for all workflow tasks),整个txtai.workflow.task包在 task/init.py 中导出了一系列内置任务:ConsoleTaskExportTaskFileTaskImageTaskRetrieveTaskServiceTaskStorageTaskStreamTaskTemplateTaskUrlTaskWorkflowTask以及别名ExtractorTask(即RagTask),它们全部继承自这个基类。

最简单的 Task

Task(lambda x: [y * 2 for y in x])

上面的 Task 对输入的所有元素执行传入的 lambda 函数——把每个元素乘以 2。注意这里的 action 接收的是一整批元素(一个列表),而不是单个元素,这与后续"one-to-many"变换、批处理的设计是一致的。

Task 与 Pipeline 的组合

由于 Pipeline 本身也是可调用对象,Task 可以直接把 Pipeline 当作 action 使用。下面这个例子对每个输入元素做摘要:

summary = Summary() Task(summary)

这种"Task 包装 Pipeline"的模式是 txtai 工作流中最常见的用法。Task 负责批处理、并发调度、结果合并等横切逻辑,Pipeline(如SummaryTranslationTranscriptionLabels等)负责真正的模型推理。可以对照 Workflow 文档 中的完整示例:FileTask(transcribe, r"\.wav$")把转写 Pipeline 包装成文件任务,随后Task(lambda x: translate(x, "fr"))再对转写出的文本执行翻译动作——这正是"tasks encapsulate processing"的直观体现。

Task 独立运行与加入 Workflow

Task 可以独立运行,但与 Workflow 配合时效果最佳,因为工作流带来了大规模流式批处理能力:

summary = Summary() task = Task(summary) task(["Very long text here"]) workflow = Workflow([task]) list(workflow(["Very long text here"]))

两种调用方式都能得到结果,但区别在于:直接调用task(...)是一次性同步处理;而workflow(...)返回一个生成器(generator),只有在被迭代消费时才真正执行,并按照 Workflow 的批处理逻辑 以batch(默认 100)为单位分块处理数据。__call__内部会先创建Execute并发执行器、运行各任务初始化方法(initialize),再逐批process,最后运行收尾方法(finalize)。这也是 Workflow 文档 反复强调的:"Since workflows run as generators, output must be consumed for execution to occur."

通过配置创建 Task

Task 也可以用配置方式创建,作为 Workflow 的一部分:

workflow: name: tasks: - action: summary

YAML 中的action既可以是内置 Pipeline 名称,也可以带task字段指定任务类型。由 TaskFactory.create 的实现可见其背后机制:配置中的args会被转成偏函数(Partial)注入 action,普通任务名会被拼接成txtai.workflow.task.<Name>Task的类路径再通过Resolver解析实例化;action支持单个或列表,列表时会按位置逐一对应args。这也解释了为何配置化的 Task 可以像下面这样传参:

workflow: index: tasks: - action: translation args: ["fr"]

Task 的完整参数体系

Task.init定义了完整的构造参数,下表逐一说明:

参数默认值说明
actionNone对每个数据元素执行的动作,可以是单个可调用对象或可调用对象列表
selectNone用于筛选待处理数据的过滤器(正则表达式);不匹配的元素原样透传
unpackTrue是否从(id, data, tag)三元组中解包出data再处理
columnNone当元素是元组时选取的列索引,默认处理全部;支持整数或{动作索引: 列索引}映射
merge"hstack"多动作输出的合并方式,可选hstack/vstack/concat/None
initializeNone处理开始前执行的动作
finalizeNone处理结束后执行的动作
concurrencyNone并发方式:"thread"线程并发、"process"进程并发,默认顺序执行
onetomanyTrue是否启用一对多数据变换(单个输入产生多个输出)
kwargs额外关键字参数,用于自定义 Task 子类的register方法

几个关键参数在源码中的行为值得展开:

  • select过滤accept方法(base.py)用re.search(self.select, element.lower())判断元素是否匹配。在filteredrun(base.py)中,每个输入元素会被打上唯一进程号,不匹配的元素原样透传,只对匹配子集执行 action——这正是 Workflow 文档 配置示例中FileTask(transcribe, r"\.wav$")只处理 wav 文件的底层实现。
  • unpack解包upack/pack(base.py)处理(id, data, tag)格式元素。解包后只对data执行动作,处理完毕后再把新数据装回元组对应位置;当解包后新数据本身就是元组且非 hstack 合并时,直接采用新元组。
  • column列提取extract方法(base.py)在元素是元组时按索引取值;column为字典时按"动作索引 → 列索引"为每个动作单独指定列。
  • onetomany一对多single(base.py)会把列表型输出包装成OneToMany对象(base.py),run/filteredpack再将其展开,实现"一个输入 → 多个输出"的流式变换。
  • initialize/finalize:分别对应 Workflow.initialize / Workflow.finalize 与 base.py 中在整批处理前后对每个任务调用的钩子,常用于打开/关闭连接、加载/释放资源。

另外注意:基类要求子类若需要额外参数,必须实现register(**kwargs)方法,否则会抛出TypeError(base.py)。TemplateTask就是一个范例,它通过register(template, rules, strict)接收模板参数(template.py)。

多动作任务的并发执行

默认情况下,多动作任务(action 为列表)按顺序执行。虽然 txtai 在多层级已经内置了并行能力——例如 GPU 模型会自动最大化 GPU 利用率、CPU 模式下也有并发——但仍存在需要"任务动作级并发"的场景,比如:

  • 系统有多个 GPU,希望多个动作并行占用不同 GPU;
  • 任务运行的是外部顺序代码;
  • 任务包含大量 I/O 操作。

此时可以把多动作任务切换为多线程多进程运行,二者的取舍如下:

  • multithreading(多线程)——没有创建独立进程和 pickle 数据的开销,但由于 GIL 的存在,Python 同一时刻只能执行一个线程,因此对 CPU 密集型动作没有帮助;非常适合 I/O 密集型动作和 GPU 动作。
  • multiprocessing(多进程)——创建独立子进程,数据通过 pickle 传递,每个进程独立运行,可以充分占满所有 CPU 核;非常适合 CPU 密集型动作。

(关于多进程的更多细节可参考 Python 官方 multiprocessing 文档。)

并发执行的底层机制

concurrency参数真正的执行入口在 Execute.run:当method非空且动作数大于 1 时,调用pool.starmap(action, inputs)对分发到线程池(ThreadPool)或进程池(Pool);否则退化为顺序的列表推导。进程池使用torch.multiprocessing.get_context("spawn")创建(execute.py),这一步会注册 PyTorch 的共享内存序列化以支持 CUDA 张量跨进程传递。整个池的生命周期由 Execute 上下文管理器 管理,Workflow 每次调用都会with Execute(self.workers)创建并最终close关闭线程/进程池。而workers的默认数量在 Workflow.init中被设为"任务中最大动作数"。

并发效果验证

测试用例 testConcurrentWorkflow 用两个Nop动作分别以"thread""process"和不合法值运行工作流,结果都得到[(2, 2), (4, 4)]——即每个输入元素被两个动作并行处理、按 hstack 合并为元组输出;非法并发值则安全回退到顺序执行。这印证了concurrency的容错设计与"并发不改变语义、只改变执行方式"的定位。

多动作任务的输出合并

多动作任务会为输入数据并行产生多路输出,Task 提供了三种合并方式(merge参数),再加上None(不合并),覆盖了几乎所有编排需求。合并逻辑集中在 postprocess:单动作任务直接single处理;多动作任务按merge分发到对应方法。

hstack:列式合并(默认)

hstack 按合并,每个输出行是各动作输出值组成的元组,是一对一变换:

Inputs: [a, b, c] Outputs => [[a1, b1, c1], [a2, b2, c2]] Column Merge => [(a1, a2), (b1, b2), (c1, c2)]

实现上,纯 Python 路径用list(zip(*outputs))完成;若所有输出都是 numpy 数组或 torch 张量,则分别走np.stack(outputs, axis=1)/torch.stack(outputs, axis=1)的原生向量化路径。testConcurrentWorkflow 中(2, 2)的结果正是 hstack 的体现。

vstack:行式合并

vstack 按合并,返回"列表的列表",会被解释为一对多变换

Inputs: [a, b, c] Outputs => [[a1, b1, c1], [a2, b2, c2]] Row Merge => [[a1, a2], [b1, b2], [c1, c2]] = [a1, a2, b1, b2, c1, c2]

即每个输入行展开成多个输出元素(每个输出仍是合并后的列表),供下游任务逐一处理。numpy / torch 路径分别用np.concatenate(np.stack(outputs, axis=1))torch.cat(...)实现。

concat:拼接为字符串

concat 先做列式合并,再把每行各动作输出用". "连接成一个字符串:

Inputs: [a, b, c] Outputs => [[a1, b1, c1], [a2, b2, c2]] Concat Merge => [(a1, a2), (b1, b2), (c1, c2)] => ["a1. a2", "b1. b2", "c1. c2"]

实现为[". ".join([str(y) for y in x if y]) for x in self.hstack(outputs)],空值会被过滤。适合把多路推理结果拼成一句可读文本(例如关键词 + 翻译结果)。

merge=None:保持多路输出

merge设为None时,postprocess直接返回未经合并的outputs列表,下游可以自行处理每一路结果。

提取 Task 输出的列

采用列式(column-wise)合并后,每个输出行是各动作输出值组成的元组。这个元组可以继续喂给下游任务,而下游任务可以通过column参数让不同动作分别处理不同元素——这正是构建复杂工作流图的关键构造。

先看一个简单示例:

workflow = Workflow([Task(lambda x: [y * 3 for y in x], unpack=False, column=0)]) list(workflow([(2, 8)]))

对于输入元组(2, 8),工作流只选取第一个元素2,对它执行乘以 3 的动作,输出6。注意这里设置了unpack=False,否则默认的unpack=True会把元组当作(id, data, tag)解包出data(即第二个元素)。

再看多动作、按列分工的版本:

workflow = Workflow([Task([lambda x: [y * 3 for y in x], lambda x: [y - 1 for y in x]], unpack=False, column={0:0, 1:1})]) list(workflow([(2, 8)]))

输入(2, 8)时,column={0:0, 1:1}表示:动作 0 取第 0 列(2)乘 3 得 6,动作 1 取第 1 列(8)减 1 得 7,再经 hstack 合并输出(6, 7)每个输入列被独立的动作处理,输出元组又可以被下一个任务继续按列拆分——正如文档所说:This simple construct can help build extremely powerful workflow graphs!

从源码看,column的字典语义在 execute 方法 中实现:index = self.column[x] if isinstance(self.column, dict) else self.column,为每个动作取出对应列索引后再extract。测试用例 testExtractWorkflow 验证了三种输入形态:普通元组(0, 1)取第 0 列返回0(0, (1, 2), None)这种(id, data, tag)结构在unpack=False时会把内层元组按列替换得到(0, 1, None);非元组输入则原样透传。

与 TemplateTask 的联动

列提取常与模板任务配合使用:TemplateTask.prepare(template.py)会把元组输入映射为arg0arg1… 的模板参数,把字典输入映射为命名参数,普通输入则作为{text}注入模板。借助column拆列 + 模板 + LLM 动作,可以在 YAML 中声明式地搭建"取第 N 列 → 套提示词 → 调 LLM"的链式工作流,相关声明式示例见 Workflow 文档 的 LLM workflow example。

完整实战:从音频到索引的混合任务流

把上述概念串起来,参考 Workflow 文档 中的完整示例,一个典型的混合任务流是:转写音频 → 翻译文本 → 写入向量索引。Python 写法如下:

from txtai import Embeddings from txtai.pipeline import Transcription, Translation from txtai.workflow import FileTask, Task, Workflow embeddings = Embeddings({ "path": "sentence-transformers/paraphrase-MiniLM-L3-v2", "content": True }) transcribe = Transcription() translate = Translation() tasks = [ FileTask(transcribe, r"\.wav$"), # 只处理 wav 文件(select 过滤) Task(lambda x: translate(x, "fr")) # 把转写结果翻译成法语 ] data = [ "US_tops_5_million.wav", "Canadas_last_fully.wav", "Beijing_mobilises.wav", "The_National_Park.wav", "Maine_man_wins_1_mil.wav", "Make_huge_profits.wav" ] workflow = Workflow(tasks) embeddings.index((uid, text, None) for uid, text in enumerate(workflow(data))) embeddings.search("wildlife", 1)

这里FileTask(transcribe, r"\.wav$")select正则只处理 wav 文件(不匹配的元素透传),Task(lambda x: translate(x, "fr"))用 lambda 包装 Translation Pipeline。整个工作流以生成器形式消费,边流式处理边喂给embeddings.index

对应的 YAML 声明式写法:

writable: true embeddings: path: sentence-transformers/paraphrase-MiniLM-L3-v2 content: true transcription: translation: workflow: index: tasks: - action: transcription select: "\\.wav$" task: file - action: translation args: ["fr"] - action: index
from txtai import Application app = Application("workflow.yml") list(app.workflow("index", [ "US_tops_5_million.wav", "Canadas_last_fully.wav", "Beijing_mobilises.wav", "The_National_Park.wav", "Maine_man_wins_1_mil.wav", "Make_huge_profits.wav" ])) app.search("wildlife")

YAML 中task: file指定FileTaskselect传入转义后的正则,args: ["fr"]通过 TaskFactory 的Partial机制注入翻译目标语言。这说明:Task 的参数体系在 Python 与配置两种形态下完全等价

更多内置 Task 一览

Task是基类,txtai 还提供了面向具体场景的内置任务(均在 task/init.py 导出,每个文件都有对应的独立文档):

任务类用途文档
ConsoleTask把数据打印到控制台docs/workflow/task/console.md
ExportTask导出为 CSV / Excel 等格式docs/workflow/task/export.md
FileTask从文件读取内容docs/workflow/task/file.md
ImageTask图像相关处理docs/workflow/task/image.md
RetrieveTask向量检索/数据库查询docs/workflow/task/retrieve.md
ServiceTask调用外部服务docs/workflow/task/service.md
StorageTask云存储读写docs/workflow/task/storage.md
StreamTask流式处理docs/workflow/task/stream.md
TemplateTask模板生成 / LLM 提示词构造docs/workflow/task/template.md
UrlTaskURL 内容抓取docs/workflow/task/url.md
WorkflowTask嵌套工作流(工作流套工作流)docs/workflow/task/workflow.md

WorkflowTask为例,它可以让你创建"工作流的工作流":Workflow([WorkflowTask(otherworkflow)])直接把另一个工作流当作任务嵌套进来(workflow.md)。这些任务全部通过TaskFactoryget方法按名称解析实例化,并继承上述全部参数能力,因此本文讲透的参数体系、并发与合并机制,对所有内置任务同样适用

小结

Task 是 txtai 工作流的最小编排单元:用action挂载可调用对象(函数或 Pipeline),用select/column/unpack控制数据选取与解包,用concurrency在多 GPU、I/O 密集或 CPU 密集场景下切换线程/进程并发,用merge(hstack/vstack/concat/None)决定多路输出的合并形态,再用column字典在下游继续按列拆分,从而组合出任意复杂的处理图。理解 Task,就等于掌握了整个 txtai 工作流体系的"砖块",再配合 Workflow 文档 与 调度文档,即可搭建完整的流式 AI 处理管线。

【免费下载链接】txtai💡 All-in-one AI framework for semantic search, LLM orchestration and language model workflows项目地址: https://gitcode.com/GitHub_Trending/tx/txtai

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询