跳转至

脚本开发 / 线程池 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
import requests

def fetch(url):
    resp = requests.get(url, timeout=10)
    resp.raise_for_status()
    return resp.text

@DFF.API('分批请求')
def batch_fetch(first_urls, second_urls):
    if not first_urls and not second_urls:
        return []

    DFF.THREAD.set_pool_size(5)
    for url in first_urls:
        DFF.THREAD.submit(fetch, url)

    # 等待第一批任务全部完成,再创建大小为 10 的线程池
    DFF.THREAD.set_pool_size(10)
    for url in second_urls:
        DFF.THREAD.submit(fetch, url)

    return [
        {'key': result.key, 'error': repr(result.error)}
        if result.error
        else {'key': result.key, 'value': result.value}
        for result in DFF.THREAD.get_all_results()
    ]

结果读取行为

  • 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 密集型任务通常不会获得明显提速。