
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 dont 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 中导出了一系列内置任务ConsoleTask、ExportTask、FileTask、ImageTask、RetrieveTask、ServiceTask、StorageTask、StreamTask、TemplateTask、UrlTask、WorkflowTask以及别名ExtractorTask即RagTask它们全部继承自这个基类。最简单的 TaskTask(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如Summary、Translation、Transcription、Labels等负责真正的模型推理。可以对照 Workflow 文档 中的完整示例FileTask(transcribe, r\.wav$)把转写 Pipeline 包装成文件任务随后Task(lambda x: translate(x, fr))再对转写出的文本执行翻译动作——这正是tasks encapsulate processing的直观体现。Task 独立运行与加入 WorkflowTask 可以独立运行但与 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.通过配置创建 TaskTask 也可以用配置方式创建作为 Workflow 的一部分workflow: name: tasks: - action: summaryYAML 中的action既可以是内置 Pipeline 名称也可以带task字段指定任务类型。由 TaskFactory.create 的实现可见其背后机制配置中的args会被转成偏函数Partial注入 action普通任务名会被拼接成txtai.workflow.task.NameTask的类路径再通过Resolver解析实例化action支持单个或列表列表时会按位置逐一对应args。这也解释了为何配置化的 Task 可以像下面这样传参workflow: index: tasks: - action: translation args: [fr]Task 的完整参数体系Task.init定义了完整的构造参数下表逐一说明参数默认值说明actionNone对每个数据元素执行的动作可以是单个可调用对象或可调用对象列表selectNone用于筛选待处理数据的过滤器正则表达式不匹配的元素原样透传unpackTrue是否从(id, data, tag)三元组中解包出data再处理columnNone当元素是元组时选取的列索引默认处理全部支持整数或{动作索引: 列索引}映射mergehstack多动作输出的合并方式可选hstack/vstack/concat/NoneinitializeNone处理开始前执行的动作finalizeNone处理结束后执行的动作concurrencyNone并发方式thread线程并发、process进程并发默认顺序执行onetomanyTrue是否启用一对多数据变换单个输入产生多个输出kwargs—额外关键字参数用于自定义 Task 子类的register方法几个关键参数在源码中的行为值得展开select过滤accept方法base.py用re.search(self.select, element.lower())判断元素是否匹配。在filteredrunbase.py中每个输入元素会被打上唯一进程号不匹配的元素原样透传只对匹配子集执行 action——这正是 Workflow 文档 配置示例中FileTask(transcribe, r\.wav$)只处理 wav 文件的底层实现。unpack解包upack/packbase.py处理(id, data, tag)格式元素。解包后只对data执行动作处理完毕后再把新数据装回元组对应位置当解包后新数据本身就是元组且非 hstack 合并时直接采用新元组。column列提取extract方法base.py在元素是元组时按索引取值column为字典时按动作索引 → 列索引为每个动作单独指定列。onetomany一对多singlebase.py会把列表型输出包装成OneToMany对象base.pyrun/filteredpack再将其展开实现一个输入 → 多个输出的流式变换。initialize/finalize分别对应 Workflow.initialize / Workflow.finalize 与 base.py 中在整批处理前后对每个任务调用的钩子常用于打开/关闭连接、加载/释放资源。另外注意基类要求子类若需要额外参数必须实现register(**kwargs)方法否则会抛出TypeErrorbase.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, axis1)/torch.stack(outputs, axis1)的原生向量化路径。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, axis1))和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)]空值会被过滤。适合把多路推理结果拼成一句可读文本例如关键词 翻译结果。mergeNone保持多路输出当merge设为None时postprocess直接返回未经合并的outputs列表下游可以自行处理每一路结果。提取 Task 输出的列采用列式column-wise合并后每个输出行是各动作输出值组成的元组。这个元组可以继续喂给下游任务而下游任务可以通过column参数让不同动作分别处理不同元素——这正是构建复杂工作流图的关键构造。先看一个简单示例workflow Workflow([Task(lambda x: [y * 3 for y in x], unpackFalse, column0)]) list(workflow([(2, 8)]))对于输入元组(2, 8)工作流只选取第一个元素2对它执行乘以 3 的动作输出6。注意这里设置了unpackFalse否则默认的unpackTrue会把元组当作(id, data, tag)解包出data即第二个元素。再看多动作、按列分工的版本workflow Workflow([Task([lambda x: [y * 3 for y in x], lambda x: [y - 1 for y in x]], unpackFalse, 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)结构在unpackFalse时会把内层元组按列替换得到(0, 1, None)非元组输入则原样透传。与 TemplateTask 的联动列提取常与模板任务配合使用TemplateTask.preparetemplate.py会把元组输入映射为arg0、arg1… 的模板参数把字典输入映射为命名参数普通输入则作为{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: indexfrom 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.mdExportTask导出为 CSV / Excel 等格式docs/workflow/task/export.mdFileTask从文件读取内容docs/workflow/task/file.mdImageTask图像相关处理docs/workflow/task/image.mdRetrieveTask向量检索/数据库查询docs/workflow/task/retrieve.mdServiceTask调用外部服务docs/workflow/task/service.mdStorageTask云存储读写docs/workflow/task/storage.mdStreamTask流式处理docs/workflow/task/stream.mdTemplateTask模板生成 / LLM 提示词构造docs/workflow/task/template.mdUrlTaskURL 内容抓取docs/workflow/task/url.mdWorkflowTask嵌套工作流工作流套工作流docs/workflow/task/workflow.md以WorkflowTask为例它可以让你创建工作流的工作流Workflow([WorkflowTask(otherworkflow)])直接把另一个工作流当作任务嵌套进来workflow.md。这些任务全部通过TaskFactory的get方法按名称解析实例化并继承上述全部参数能力因此本文讲透的参数体系、并发与合并机制对所有内置任务同样适用。小结Task 是 txtai 工作流的最小编排单元用action挂载可调用对象函数或 Pipeline用select/column/unpack控制数据选取与解包用concurrency在多 GPU、I/O 密集或 CPU 密集场景下切换线程/进程并发用mergehstack/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),仅供参考