【玩转daft】udf的几种使用方式

发布时间:2026/9/16 14:11:12

【玩转daft】udf的几种使用方式 daft.udf 已经0.7.0正式标记 deprecated。daft提供非常灵活的函数定义形式。1对1 row-rise1 row in - 1 value out 很多算子的组织形式importdaftdaft.funcdefadd(a:int,b:int)-int:returnab dfdaft.from_pydict({x:[1,2],y:[10,20]})dfdf.with_column(z,add(df[x],df[y]))1对多1 row in - N rows out可以这么写fromtypingimportIteratordaft.funcdefsplit_into_sentences(text:str)-Iterator[str]:importre sentencesre.split(r(?[.!?])\s,text.strip())forsentenceinsentences:ifsentence:yieldsentence输入输出这里Daft 会自动把原来的 ticket_id 复制到每个生成出来的新行不需要手动 explode()。这里的关键是 Iterator[T] 和 yield。返回类型是 Iterator[str]。这告诉 Daft这个函数不是返回一个普通字符串而是会连续 yield 多个字符串。比如yieldchunk 1yieldchunk 2yieldchunk 3Daft 就知道每个 yield 出来的值都应该变成一行。传统的做法 list explode()daft.funcdefsplit_into_sentences(text:str)-list[str]:return[A.,B.,C.]dfdf.with_column(sentences,split_into_sentences(df[body]))dfdf.explode(sentences)不用先把所有结果收集成一个大 list再 explode。对于长文档、音频切片、日志拆分、视频帧采样这种场景Generator 更自然也更省内存。访问外部接口 使用异步方式访问外部的时候使用异步的方式。适合 I/O比如 HTTP API、对象存储、小文件下载、远程 embedding 服务。daft.func(max_concurrency10)asyncdeffetch_url(url:str)-str:importaiohttpasyncwithaiohttp.ClientSession()assession:asyncwithsession.get(url)asresponse:returnawaitresponse.text()dfdf.with_column(html,fetch_url(df[url]))max_concurrency10 表示限制并发请求数避免打爆 API。避免内存/网络压力过大。异步函数asyncdeffetch(url):...调用后不是马上把结果算出来而是返回一个 coroutine需要被事件循环调度执行。在 Daft 里async def 的意义是这个 UDF 可以并发执行多行而不是一行一行阻塞等待。比如 100 个 URL 请求普通同步函数请求 1 完成再请求 2再请求 3async 函数可以同时发多个请求谁先返回就先处理谁。有状态的访问有昂贵初始化时用 daft.cls。比如加载模型、初始化 tokenizer、创建数据库连接。daft.clsclassTextClassifier:def__init__(self,model_path:str):self.modelload_model(model_path)def__call__(self,text:str)-str:returnself.model.predict(text)classifierTextClassifier(model.pkl)dfdf.with_column(label,classifier(df[text]),)不会立刻真的加载模型。Daft 执行任务时会在 worker 上初始化实例并复用它处理多行。daft.cls 里可以有多个方法daft.clsclassTextProcessor:def__init__(self,prefix:str):self.prefixprefixdef__call__(self,text:str)-str:returnself.prefixtextdeflowercase(self,text:str)-str:returntext.lower()deflength(self,text:str)-int:returnlen(text)processorTextProcessor( )dfdf.select(processor(df[text]).alias(prefixed),processor.lowercase(df[text]).alias(lower),processor.length(df[text]).alias(length),)call可以直接processor(df[text])普通方法要processor.lowercase(df[text])在 daft.cls 里如果某个方法需要指定返回类型用 daft.method。fromdaftimportDataTypedaft.clsclassTextProcessor:daft.method(return_dtypeDataType.list(DataType.string()))defsplit_words(self,text:str):returntext.split()或者返回 structdaft.clsclassAnalyzer:daft.method(return_dtypedaft.DataType.struct({word_count:daft.DataType.int64(),char_count:daft.DataType.int64(),}),unnestTrue,)defanalyze(self,text:str):return{word_count:len(text.split()),char_count:len(text),}调用dfdf.select(Analyzer().analyze(df[text]))如果 unnestTruestruct 会展开成多列word_count /char_count批处理daft.func.batch它和普通 daft.func 的区别是普通 UDF一行一行处理Batch UDF一批一批处理daft.func.batch(return_dtypedaft.DataType.int64())defword_count_batch(texts:daft.Series)-list:Count words in each text -- operating on the entire batch at once.return[len(text.split())fortextintexts.to_pylist()]这里函数收到的不是一个 str而是一批文本texts: daft.Series可以理解成[hello world,this is a ticket,daft udf batch example]然后函数返回同样长度的结果[2,4,4]也就是Series[str] - list[int]主要是为了减少 Python 调用开销并且方便用向量化库。普通 UDF每一行调用一次 Python 函数。比如 100 万行可能要调用 100 万次。Batch UDF 每一批调用一次 Python 函数 。比如每批 1024 行100 万行大约调用 1000 次。调用次数少很多。适合 NumPydaft.func.batch(return_dtypedaft.DataType.float64())defnormalize(values:daft.Series)-list:importnumpyasnp arrnp.array(values.to_pylist())arr(arr-arr.mean())/arr.std()returnarr.tolist()适合 pandasdaft.func.batch(return_dtypedaft.DataType.string())defclean_texts(texts:daft.Series)-list:importpandasaspd spd.Series(texts.to_pylist())ss.str.lower().str.strip()returns.tolist()适合批量 API比如 embedding API 通常支持一次传多个文本daft.func.batch(return_dtypedaft.DataType.embedding(daft.DataType.float32(),1536))defembed_texts(texts:daft.Series)-list:client...responseclient.embeddings.create(inputtexts.to_pylist(),modeltext-embedding-3-small,)return[item.embeddingforiteminresponse.data]这比每行单独请求一次 API 高效很多。参考https://docs.daft.ai/en/stable/examples/udf-patterns/#pattern-2-generator-one-input-becomes-many-rows
延伸阅读

更多相关文章

2026/9/16 14:06:10

PHP原生CC防护系统:请求指纹+动态限流+验证码闭环

简介:这是一套面向PHP开发者与Web安全初学者的轻量级CC攻击防护实践源码,聚焦于解决PHP网站在高并发场景下易遭模拟请求式DDoS(即CC攻击)导致服务瘫痪的问题。资源共14个文件,含7个核心PHP脚本(如anti_ddos…

2026/9/16 14:06:10

C51单片机实时频谱分析:滑动DFT定点实现与Keil工程落地

简介:本资源是面向嵌入式开发工程师与数字信号处理学习者的C51平台滑动DFT(Sliding DFT)完整实现方案,聚焦实时频谱分析场景下的轻量级算法落地。针对C51单片机内存与算力受限的特点,资源提供可直接移植的优化代码、原…

2026/9/16 15:01:31

51单片机+Proteus三路抢答器仿真设计与调试

简介:本资源是一套基于51单片机开发的八路智能抢答器完整设计资料,面向电子类专业初学者、课程设计学生及单片机实践爱好者,解决课堂互动、知识竞赛等场景下的实时抢答与计分需求。压缩包共16个文件,涵盖Proteus仿真工程&#xff…

2026/9/16 15:01:31

FPGA出租车计费器设计:硬件级实时计程计时与LED动态扫描实现

简介:本资源是一套面向数字电路与FPGA初学者及课程设计学生的完整出租车计费系统实现方案,聚焦VHDL语言开发与EDA工程实践,解决嵌入式计费逻辑建模、多模块协同仿真与FPGA硬件部署等核心问题。压缩包共245个文件,涵盖57个.cdb与57…

2026/9/16 15:01:31

agent-skills:面向智能体的前端工程化技能抽象范式

1. “agent-skills”不是库名,而是工程级能力抽象范式你搜“agent-skills”,首页跳出来的全是 TypeScript、Node、Nx 相关的开发问题——npm 报错、PowerShell 脚本被禁、nvm 切换失败、TS 类型报错、Nx 插件加载异常……但没人解释“agent-skills”本身…

2026/9/16 12:52:37

拯救者Y7000黑屏故障排查与维修实战指南

1. 项目概述:一台黑屏的拯救者Y7000,到底卡在哪一步? 联想拯救者Y7000系列笔记本,从2018年第一代搭载i5-8300H开始,到后来的i7-9750H、i7-10750H、i5-11400H,再到2023年款的R7-7840HS,它始终是学…

2026/9/16 0:04:09

PHP源码部署实战:从环境配置到运行情侣游戏全攻略

简介:这是一套面向情侣互动场景的PHP完整源码,集成情侣飞行棋、真心话大冒险、情趣骰子等玩法,并内置完整分销制度,可自定义多种返佣比例,源码完全开源无加密,支持微信无感自动授权登录与第三方授权&#x…

2026/9/15 14:22:53

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/15 21:31:11

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/15 11:42:23

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码