脚本开发 / 线程池 DFF.THREAD
DFF.THREAD 用于在单个 Task 内并发执行 IO 密集型函数,例如批量 HTTP 请求。线程池由 DataFlux Func 管理。
API
| 方法 / 属性 | 说明 |
|---|---|
pool_size |
已显式配置的线程池大小;尚未设置时为 None |
set_pool_size(pool_size) |
设置线程池大小;已有线程池时等待全部任务完成后重建 |
submit(fn, *args, **kwargs) |
提交函数并返回结果 Key |
get_result(key, wait=True) |
获取指定结果 |
get_all_results(wait=True) |
获取全部结果 |
pop_result(wait=True) |
弹出一个已经完成的结果;弹出后不会再被其他方法返回 |
is_all_finished |
是否全部执行完成 |
wait_all_finished() |
等待全部执行完成 |
结果对象类型为 DFFThreadResult,包含以下属性:
| 属性 | 说明 |
|---|---|
key |
submit(...) 返回的结果 Key |
value |
函数返回值;执行失败时通常为 None |
error |
函数抛出的异常;执行成功时为 None |
调整线程池大小
set_pool_size(...) 只能在 Task 主线程中调用,并会按照在配置文件 config.yaml 中配置的最大值进行校验。传入无效大小时会抛出异常,已存在的线程池保持不变。
首次 submit(...) 前调用时,只会记录设置的大小,线程池在提交任务时创建。未显式设置时,创建线程池会使用运行时默认值。
自 8.1.15 起,已有线程池时再次调用 set_pool_size(...),会等待全部已提交任务(包括排队任务)完成,关闭旧线程池并按指定大小创建新线程池。即使指定的大小与当前相同,也会等待并重建;调用返回后再提交下一批任务。
调整大小会保留已有的结果 Key、返回值和异常,尚未收集的结果仍可使用原有结果读取方法获取。
示例
| 分两批执行 HTTP 请求 | |
|---|---|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 | |
结果读取行为
get_result(key, wait=False)和pop_result(wait=False)在没有匹配的就绪结果时返回None。- 只有至少提交过一个函数后才能调用
get_all_results(...)。 get_all_results(...)按提交 Key 的顺序返回结果,不按完成顺序返回。- 需要边完成边处理时,可以调用
pop_result(wait=not DFF.THREAD.is_all_finished)。 pop_result(...)适合互相独立的任务;get_all_results(...)适合全部完成后统一处理。
Warning
每个 result.error 都必须检查,否则线程异常会被忽略。即使结果尚未收集,Task 也会等待已提交的线程工作停止后才结束,因此每个线程操作都必须设置明确的超时或其他执行边界。
线程池适合 IO 密集型工作;CPU 密集型任务通常不会获得明显提速。