如何让Faust Agent处理完一批数据后再取下一批(同步运行)?
实现Faust Agent串行批次处理的方案
没问题,这个需求很容易实现——你只需要把批次处理的逻辑放到一个无限循环里,让Agent在处理完当前批次后,主动去取下一批数据就行。
原来的代码只执行了一次take操作,处理完那5000条(或5秒内的不足量)就停止了。改成循环之后,就能持续按照「取一批→处理完→再取下一批」的串行逻辑运行。
具体实现代码
@app.agent() async def process(stream): while True: # 等待收集最多5000条记录,超时时间保持5秒 batch = await stream.take(5000, within=5) # 避免处理空批次(比如5秒内没有任何数据的情况) if not batch: continue # 同步处理批次内的每条数据(如果process是异步函数记得加await) for value in batch: process(value) # 若process为异步,改为 await process(value) # 或者如果你的处理逻辑支持批量操作,直接传入整个批次更高效 # process_batch(batch)
关键逻辑说明
while True:让Agent持续运行,不断处理新的批次await stream.take(...):这一步是阻塞式的,会等待直到收集到5000条记录,或者5秒超时(取先满足的条件),只有拿到批次后才会进入处理环节- 处理完当前批次的所有数据后,循环才会回到开头,再次调用
take取下一批,完美实现「处理完上一批再取下一批」的同步效果
如果你的场景需要必须凑够5000条才处理(不接受超时不足量的情况),只需要去掉within=5参数即可:
batch = await stream.take(5000)
这样Agent会一直等待,直到收集满5000条记录才开始处理。
内容的提问来源于stack exchange,提问作者Rohit
相关产品推荐
相关产品推荐

