Lewati ke isi

Pengembangan Skrip / Langganan Konektor

Sebagian Konektor mendukung berlangganan pesan, DataFlux Func menyediakan cara terpadu untuk berlangganan.

Pada versi sebelumnya, disebut 'Data Source', sekarang diubah menjadi 'Konektor'

1. Pendahuluan

Karena mekanisme eksekusi Skrip DataFlux Func, setelah fungsi dimulai, fungsi tetap harus berakhir, dan tidak diizinkan berjalan tanpa batas.

Oleh karena itu, untuk penanganan proses jangka panjang seperti langganan, tidak didukung untuk secara langsung menulis konsumen yang berjalan lama di DataFlux Func untuk mewujudkannya.

Sebaliknya, perlu menentukan topik langganan di Konektor, dan program utama DataFlux Func-lah yang secara terpusat menangani langganan pesan.

Ketika program utama DataFlux Func menerima pesan, pesan akan diteruskan ke fungsi penangan pesan yang ditentukan untuk diproses, sehingga menyelesaikan penanganan langganan pesan.

2. Konektor yang Mendukung Langganan

DataFlux Func versi terbaru mendukung langganan untuk Konektor berikut:

3. Langkah Operasi

Untuk berlangganan pesan Redis dan memprosesnya, langkah-langkah operasi spesifiknya adalah sebagai berikut:

3.1 Menulis Fungsi Penangan Pesan

Fungsi penangan pesan memiliki bentuk fungsi yang tetap, seperti:

Python
1
2
3
4
@DFF.API('Message Handler')
def message_handler(topic, message):
    print('topic', topic)     # topik
    print('message', message) # konten

Setelah menyelesaikan Skrip, simpan dan publikasikan:

message-handler-func.png

3.2 Membuat dan Mengonfigurasi Konektor

Selain mengisi konfigurasi dasar, Anda juga perlu mengisi topik langganan (Topic) serta fungsi penangan yang sesuai.

Fungsi penangan adalah fungsi message_handler(...) yang telah ditulis di atas:

sub-config.png

3.3 Menerbitkan Pesan dan Mengonfirmasi Penanganan Pesan

Setelah penerbit menerbitkan pesan seperti berikut, pesan akan diteruskan melalui program utama DataFlux Func ke fungsi message_handler(...) di atas untuk diproses:

Di halaman konfigurasi Konektor, di bawah topik yang sesuai akan muncul informasi konsumsi terbaru:

sub-info.png

Setelah diklik, Anda dapat melihat informasi tugas yang lebih rinci:

sub-info-detail.png

Dengan demikian, dapat dikonfirmasi bahwa pesan yang diterbitkan di Redis benar-benar diterima dan diproses oleh fungsi message_handler(...).

4. Menerbitkan Pesan

Konektor yang mendukung langganan pesan umumnya juga mendukung penerbitan (publish) pesan.

Untuk penerbitan pesan, lihat API objek Konektor yang sesuai di Pengembangan Skrip / Objek Konektor DFF.CONN, umumnya dalam bentuk .publish(topic, message).

5. Catatan Tugas Fungsi Penangan Pesan

Di DataFlux Func versi terbaru, Anda dapat langsung melihat catatan tugas dari fungsi penangan pesan.

message-handler-task-info.png

Jika jumlah pesan yang dilanggani sangat besar, mencatat log penanganan setiap pesan itu sendiri juga dapat menyebabkan masalah kinerja. Anda dapat merujuk ke Penerapan dan Pemeliharaan / Metrik Sistem dan Catatan Tugas / Menonaktifkan Catatan Tugas Fungsi Lokal untuk menonaktifkan 'Catatan Tugas Fungsi Lokal', sehingga mengurangi tekanan penyimpanan MySQL

Di DataFlux Func versi lama, fungsi penangan pesan tidak mendukung kueri catatan tugas. Jika ingin mencatat kesalahan yang terjadi selama pemrosesan, Anda dapat menulis informasi terkait ke DFF.CACHE di dalam fungsi penangan pesan.

Lihat implementasi kode berikut:

Python
 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 arrow
import traceback

def message_handler_impl(topic, message):
    # Fungsi pemrosesan pesan sebenarnya berada di sini
    # Di sini diasumsikan terjadi kesalahan pembagian dengan 0
    x = 100 / 0

@DFF.API('Message Handler')
def message_handler(topic, message):
    try:
        # Panggil fungsi pemrosesan pesan yang sebenarnya
        message_handler_impl(topic, message)

    except Exception as e:
        # Dapatkan waktu saat ini
        now_str = arrow.now('Asia/Shanghai').format('YYYY-MM-DD HH:mm:ss')

        # Ambil informasi kesalahan lengkap
        error_stack = traceback.format_exc()

        # Simpan informasi kesalahan ke DFF.CACHE
        latest_error = '\n'.join([ 'Waktu:', now_str, 'Informasi kesalahan:', error_stack ])
        DFF.CACHE.set('latest_error', latest_error)

        # Kesalahan harus dilemparkan kembali
        raise

Setelah muncul masalah selama proses eksekusi, Anda dapat melihat informasi detail di 'Manajemen / Manajer Cache Fungsi'

6. Langganan Tunggal dan Multi-Langganan

Ditambahkan pada versi 1.7.31

Semua subscriber berjalan di layanan server DataFlux Func. Setelah subscriber menerima pesan, tugas eksekusi fungsi akan dibuat sesuai konfigurasi, lalu dikirim ke worker untuk dieksekusi.

Subscriber mendukung 2 cara berlangganan:

  1. Langganan tunggal: yaitu berapa pun jumlah replika server yang diaktifkan, subscriber akan selalu berjalan hanya di salah satu server tersebut
  2. Multi-Langganan: ketika beberapa replika server diaktifkan, subscriber akan berjalan di setiap server tersebut

Untuk subscriber Redis, karena tidak mendukung penanganan langganan bersama, multi-langganan hanya akan menyebabkan konsumsi pesan yang berulang, oleh karena itu saat ini hanya mendukung langganan tunggal.

Untuk subscriber MQTT dan Kafka, Anda dapat memilih langganan tunggal atau multi-langganan.

Saat MQTT menggunakan multi-langganan, Topic perlu melakukan langganan bersama dengan format $share/group_name/topic_name, jika tidak, akan menyebabkan penerimaan pesan yang sama secara berulang

Tentang layanan yang termasuk dalam DataFlux Func, serta mengaktifkan beberapa replika server, silakan lihat Penerapan dan Pemeliharaan / Arsitektur, Penskalaan, dan Pembatasan Sumber Daya

7. Meningkatkan Kecepatan Pemrosesan Langganan

Ditambahkan pada versi 1.7.31

Untuk DataFlux Func, karena subscriber berjalan di layanan server, tetapi fungsi berjalan di layanan worker, kecepatan pemrosesan langganan terdiri dari dua bagian: kecepatan penerimaan pesan langganan dan kecepatan eksekusi fungsi langganan.

Jika jumlah pesan yang dilanggani tidak banyak, tetapi pemrosesan setiap pesan cukup kompleks dan memakan waktu lama, maka perlu menambah jumlah worker untuk memastikan pesan diproses tepat waktu.

Jika jumlah pesan yang dilanggani sangat besar, maka perlu meningkatkan jumlah server dan worker secara bersamaan, serta mengatur subscriber ke mode 'Multi-Langganan'.

Tentang cara melakukan penskalaan sistem, silakan lihat Penerapan dan Pemeliharaan / Arsitektur, Penskalaan, dan Pembatasan Sumber Daya

8. Batasan Langganan

Ditambahkan pada versi 1.7.31

Ketika jumlah pesan langganan sangat besar, karena cara publish-subscribe di sisi server berbeda, maka terdapat batasan yang berbeda pula.

Redis / MQTT

Subscriber Redis dan MQTT, karena tidak dapat mengontrol kecepatan penerimaan pesan di sisi subscriber, maka pesan yang diterima akan masuk ke dalam buffer pool internal.

Ukuran default buffer pool adalah 5000, yaitu memungkinkan maksimal 5000 pesan langganan menunggu untuk diproses. Jika selama waktu tersebut ada lebih banyak pesan yang diterima, pesan yang meluap akan dibuang.

Untuk mengatasi masalah di atas, Anda dapat meningkatkan kecepatan pemrosesan pesan langganan dengan menambah jumlah worker.

Kafka

Subscriber Kafka, karena dapat mengontrol kecepatan penerimaan pesan di sisi subscriber, maka tidak memiliki buffer pool, dan juga tidak akan membuang pesan.

Ketika kecepatan pemrosesan pesan tidak dapat mengimbangi kecepatan penerbitan pesan, Anda dapat meningkatkan kecepatan pemrosesan pesan langganan dengan menambah jumlah worker.

9. Masalah yang Diketahui

Masalah yang diketahui saat ini adalah sebagai berikut.

MQTT Macet, Tidak Dapat Menerima Pesan

MQTT Broker dapat mengalami masalah macet atau tidak dapat menerima pesan saat terhubung ke EMQX.

Masalah ini pernah muncul pada EMQX edisi komunitas single-node versi awal, sedangkan versi EMQX lainnya belum diuji.

Tetapi masalah ini tidak terlihat saat terhubung ke mosquitto.

Kafka Tidak Dapat Langsung Melakukan Konsumsi Setelah Dimulai

Setelah berlangganan, Konektor Kafka mungkin berada dalam status jeda selama beberapa menit pertama, dan baru kemudian dapat melakukan konsumsi pesan secara normal.

Penyebabnya mungkin terkait dengan mekanisme Rebalance Kafka.