Lewati ke isi

Pengembangan Skrip / Kumpulan Utas DFF.THREAD

DFF.THREAD digunakan untuk mengeksekusi fungsi intensif IO secara bersamaan dalam satu Task, misalnya permintaan HTTP massal. Kumpulan utas dikelola oleh DataFlux Func.

API

Metode / Atribut Deskripsi
pool_size Ukuran kumpulan utas yang dikonfigurasi secara eksplisit; None jika belum diatur
set_pool_size(pool_size) Mengatur ukuran kumpulan utas, harus dipanggil sebelum submit(...) pertama kali
submit(fn, *args, **kwargs) Mengirimkan fungsi dan mengembalikan Key hasil
get_result(key, wait=True) Mendapatkan hasil yang ditentukan
get_all_results(wait=True) Mendapatkan semua hasil
pop_result(wait=True) Mengeluarkan satu hasil yang telah selesai; setelah dikeluarkan tidak akan dikembalikan oleh metode lain
is_all_finished Apakah semua telah selesai dieksekusi
wait_all_finished() Menunggu semua selesai dieksekusi

Tipe objek hasil adalah DFFThreadResult, dengan atribut berikut:

Atribut Deskripsi
key Key hasil yang dikembalikan oleh submit(...)
value Nilai pengembalian fungsi; biasanya None saat eksekusi gagal
error Pengecualian yang dilemparkan fungsi; None saat eksekusi berhasil

set_pool_size(...) akan divalidasi sesuai nilai maksimum yang dikonfigurasi dalam file konfigurasi config.yaml; pengaturan kembali setelah pengiriman pertama tidak akan berlaku. Jika tidak diatur secara eksplisit, pembuatan kumpulan utas akan menggunakan nilai default runtime.

Contoh

Permintaan HTTP Massal
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
import requests

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

@DFF.API('Permintaan Massal')
def batch_fetch(urls):
    DFF.THREAD.set_pool_size(10)

    for url in 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()
    ]

Perilaku Pembacaan Hasil

  • get_result(key, wait=False) dan pop_result(wait=False) mengembalikan None ketika tidak ada hasil siap yang cocok.
  • get_all_results(...) hanya dapat dipanggil setelah setidaknya satu fungsi dikirimkan.
  • get_all_results(...) mengembalikan hasil sesuai urutan Key yang dikirimkan, bukan urutan penyelesaian.
  • Jika perlu memproses hasil seiring penyelesaiannya, Anda dapat memanggil pop_result(wait=not DFF.THREAD.is_all_finished).
  • pop_result(...) cocok untuk tugas yang saling independen; get_all_results(...) cocok untuk pemrosesan terpadu setelah semuanya selesai.

Warning

Setiap result.error harus diperiksa; jika tidak, pengecualian utas akan diabaikan. Meskipun hasil belum dikumpulkan, Task akan tetap menunggu pekerjaan utas yang telah dikirim berhenti sebelum berakhir; oleh karena itu, setiap operasi utas harus menetapkan batas waktu yang jelas atau batas eksekusi lainnya.

Kumpulan utas cocok untuk pekerjaan intensif IO; tugas intensif CPU biasanya tidak akan mendapatkan percepatan yang signifikan.