ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

【玩转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
返回列表