Azure HorizonDB için pg_durable ile Dayanıklı İşlevler (Önizleme)

pg_durable, Azure HorizonDB’nin içindeki dayanıklı yürütme altyapısıdır. Postgres'ten çıkmadan uzun süre çalışan, çok adımlı SQL iş akışlarını (işlem hatlarını ekleme, ETL işleri, yapay zeka çağrıları, zamanlanmış işler, onay akışları) tanımlamanızı ve bunları Dayanıklı İşlevler gibi ayrılmış bir düzenleyiciden beklediğiniz güvenilirlik garantileriyle çalıştırmanızı sağlar.

pg_durable aynı zamanda dayanıklı yapay zeka işlem hatlarının altındaki yürütme katmanıdır. Yapay zeka işlem hatlarını kullanıyorsanız kilitlenmelerden pg_durable kurtulmalarını, hata durumunda yeniden denemelerini ve son tamamlanan adımdan devam etmelerini sağlayan şey budur.

Note

pg_durable önizleme aşamasındadır.

"Dayanıklı" ne anlama gelir?

pg_durable içindeki dayanıklı işlev, sürecin her adımında diske kalıcı olarak yazılır. Bu size düz BEGIN ... COMMIT blok veya cron işinden almayacağınız belirli bir garanti kümesi sağlar:

  • Veritabanı çökmelerine ve yeniden başlatmaya rağmen çalışmaya devam eder. Tamamlanan adımlar, sunucu yeniden açıldığında yeniden yürütülemez. Devam eden adımlar son denetim noktasından devam eder. Çalışan yeniden çevrimiçi olduğunda bekleyen adımlar çalıştırılır.
  • Uzun bekleme sürelerine dayanır. İş akışı saatlerce uyuyabilir, cron zamanlamasını bekleyebilir veya bir dış sinyali engelleyebilir ve yine de kaldığı yerden devam edebilir.
  • Arızalara dayanır. Başarısız adımlar, işlevin tamamı yeniden çalıştırılmadan otomatik olarak yeniden denenebilir.
  • Kimliği algılar. İşlev, arka plan çalışanının ayrıcalıklarıyla değil, onu başlatan kullanıcının ayrıcalıklarıyla yürütülür. Çok kiracılı iş yükleri yalıtılmış kalır.
  • SQL'den gözlemlenebilir durumda kalır. Durumu, geçmişi, çalıştırma sayısını ve çıktıları, HorizonDB’de diğer her şey için kullandığınız aynı arayüz üzerinden inceleyebilirsiniz: bir SELECT deyimi.

Dayanıklılığın otomatik olarak yapmadığı şey: idempotent olmayan dış işlemleri kendi başına güvenle yeniden denenebilir hale getirmez. Bir adım ücret alan bir dış API’yi çağırıyorsa, bu adımı idempotent olacak şekilde tasarlayın (örneğin, bir idempotency anahtarı göndererek).

pg_durable ne zaman kullanılır?

Bu şekilde çalışmanız gerektiğinde pg_durable kullanın:

  • İşin ortasında başarısız olacak kadar uzun sürer (milyonlarca satırda embedding oluşturma, çok adımlı bir ETL süreci, geriye dönük doldurma işlemi).
  • Zaten başarılı olan bölümleri yinelemeden hata durumunda yeniden denenmesi gerekir.
  • Bir zamanlamaya göre çalışması gerekir (her saat, her hafta içi 09:00).
  • Bir dış olayı (onay, web kancası, başka bir sistemden gelen sinyal) beklemesi gerekir.
  • Dallanma, katılma veya yarış ile birden çok adımı koordine eder.
  • Şu anda dış düzenleyici + Postgres veritabanı olarak uygulanıyor ve burada işin çoğu veritabanı bölümüdür.

İş yükünüz tek bir kısa işlemsel ifadeden oluşuyorsa, pg_durable öğesine ihtiyacınız yoktur. Normal INSERT / UPDATEkullanın.

Nasıl çalışır?

Dayanıklı işlev, SQL DSL ile oluşturduğunuz ve ile df.start()gönderdiğiniz adımların grafiğidir. Graf kalıcı olarak kaydedilir, ardından bir arka plan işçisi onu yürütür.

İki önemli fikir:

  • İşlev grafı ve yürütme durumu HorizonDB'nin kendisinde ve df şemalarında duroxide depolanır. Yedeklemeler, belirli bir noktaya geri yükleme ve yüksek kullanılabilirlik, iş akışı durumunuz için otomatik olarak geçerlidir. Yönetecek ayrı düzenleyici durumu yok.
  • Arka plan çalışanı, shared_preload_libraries tarafından başlatılır. CREATE EXTENSION'den sonra uzantıyı algılar ve işlevleri çalıştırmaya başlar. Veritabanı yeniden başlatılırsa çalışan, çalışan örneklere yeniden bağlanır ve bunları sürdürür.

Note

pg_durable içindeki yürütme altyapısı, Rust için Microsoft’un açık kaynaklı dayanıklı yürütme çalışma zamanı olan Duroxide üzerine kuruludur (Durable Task Framework ve Temporal’dan esinlenmiştir). Şema duroxide adı şunu yansıtır: Duroxide düzenleme geçmişini, bağıntı kimliklerini ve yeniden yürütme durumunu burada sürdürür. pg_durable’den aldığınız deterministik yeniden oynatma, ilişkili olay kimliği ve kalıcı zamanlayıcı garantileri doğrudan Duroxide’dan gelir.

pg_durable etkinleştirme

Azure HorizonDB'de pg_durable etkinleştirmek için önce bir parametre grubu yapılandırın, ardından uzantıyı her veritabanında oluşturun.

Şu kurulum makalelerini kullanın:

  1. Sunucunuz için bir parametre grubu oluşturun.
  2. shared_preload_libraries öğesini pg_durable içerecek şekilde ayarlayın.
  3. azure.extensions öğesini pg_durable içerecek şekilde ayarlayın.
  4. Parametre grubunu sunucuya uygulayın.
  5. Her hedef veritabanına bağlanın ve şu komutu çalıştırın:

Uzantıyı kullanmak istediğiniz her veritabanında oluşturun:

CREATE EXTENSION IF NOT EXISTS pg_durable;

CREATE EXTENSION şemayı df (işlev grafikleri ve izleme görünümleri) ve şemayı duroxide (yürütme durumu) sağlar. Arka plan çalışanı birkaç saniye içinde uzantıyı algılar ve işlevleri çalıştırmaya hazırdır.

İlk dayanıklı işleviniz

-- Start a one-step durable function
SELECT df.start('SELECT ''Hello, durable world!''');
-- Returns an 8-character instance ID, for example: a1b2c3d4

-- Check status
SELECT df.status('a1b2c3d4');

-- Get the result
SELECT df.result('a1b2c3d4');

Tek adımlı bir işlev bile dayanıklıdır: veritabanı df.start() sonrasında ve worker onu almadan önce yeniden başlatılsa bile işlev yine de çalışır.

Note

df.start() bir iş akışını zaman uyumsuz olarak gönderir ve hemen döndürür. Çok adımlı iş akışları için, yan etkileri doğrulamadan önce tamamlanmasını onaylamak için , df.list_instances(), df.instance_info()veya df.status() kullanındf.result().

Program modeli

Dayanıklı işlev, adımlardan, işleçlerden ve yerleşik işlevlerden oluşturulmuş bir grafiktir. Düz SQL dizeleri otomatik olarak eşlenir, bu nedenle açıkça çağırmanız df.sql() gerekmez.

Operators

Operator Meaning Example
~> Sekans - önce sola, ardından sağa çalıştır 'SELECT 1' ~> 'SELECT 2'
& Join - paralel olarak çalıştır, tümünü bekle 'SELECT 1' & 'SELECT 2'
| Yarış - paralel çalıştır, ilk kazanır fast_query | df.sleep(30)
?> !> If / else - Boolean bir koşula göre dallanma cond ?> then_branch !> else_branch
@> Döngü - sonsuza kadar yinele (ön ek işleci) @> body
|=> Adlandır - bir adımın sonucunu kaydetmek 'SELECT id FROM users LIMIT 1' |=> 'user_id'

Kullanışlı yerleşik özellikler

Function Purpose
df.sleep(seconds) N saniye boyunca duraklat. Yeniden başlatmalara karşı kalıcıdır.
df.wait_for_schedule(cron) Bir sonraki cron ifadesi eşleşmesine kadar bekleyin.
df.wait_for_signal(name, timeout) Bir dış df.signal() gelene kadar engelleyin.
df.http(url, method, body, headers, timeout) Geçici hata durumunda yeniden deneyerek HTTP çağrısını kalıcı bir etkinlik olarak yapın.
df.if(cond, then, else) Koşullu dallanma.
df.loop(body, cond) SQL koşulu doğru olsa da tekrarlayın.
df.join(a, b) / df.race(a, b) Paralel ve yarış yürütme.
df.join3(a, b, c) Üç yönlü paralel yürütme için.
df.start(body, label, database) Dayanıklı bir işlev gönderin ve örnek kimliğini döndürin.
df.cancel(id, reason) Çalışan bir örneği iptal edin.
df.status(id) / df.result(id) Sonucu inceleyin.
df.explain(input) Görselleştirme için işlev grafını işleme.

Tüm pg_durable özellikleri hakkında daha fazla bilgi edinin.

Variables

|=>, bir adımın sonucunu bir ad altında kaydeder; sonraki adımlar bu sonuca $name olarak başvurabilir.

SELECT df.start(
    'SELECT 100 AS amount' |=> 'total'
    ~> 'SELECT $total * 2 AS doubled'
);

Kullanım örnekleri

Yeniden denemelerle çok adımlı ETL

Temizleyen, yükleyen, dizinleyen ve günlüğe kaydeden günlük ETL:

SELECT df.start(
    'DELETE FROM target WHERE loaded_at < now() - INTERVAL ''1 day'''
    ~> 'INSERT INTO target SELECT * FROM staging'
    ~> 'REINDEX TABLE target'
    ~> 'INSERT INTO etl_log (job, finished_at) VALUES (''nightly'', now())',
    'nightly-etl'
);

Veritabanı DELETE ile INSERT arasında yeniden başlatılırsa, worker INSERT noktasından devam eder - DELETE öğesini yeniden çalıştırmaz.

Zamanlanmış iş (cron)

Haftanın her günü saat 09:00'da bir bakım görevi çalıştırın:

SELECT df.start(
    @> (
        df.wait_for_schedule('0 9 * * 1-5')
        ~> 'CALL refresh_materialized_views()'
    ),
    'weekday-refresh'
);

Bu işi durdurmak istiyorsanız, işlevi çalıştırabilirsiniz cancel .

SELECT df.cancel('a1b2c3d4', 'stop test cron job');

Zaman aşımı ile onay iş akışı

Harici bir onay sinyalini 24 saate kadar bekleyin, ardından onaylayın veya reddedin:

SELECT df.start(
    'SELECT order_id, total FROM orders WHERE id = 1' |=> 'order'
    ~> df.wait_for_signal('approval', 86400) |=> 'sig'
    ~> df.if(
        'SELECT NOT ($sig::jsonb->>''timed_out'')::boolean
            AND ($sig::jsonb->''data''->>''approved'')::boolean',
        'UPDATE orders SET status = ''approved'' WHERE id = $order_id',
        'UPDATE orders SET status = ''rejected'' WHERE id = $order_id'
    ),
    'order-approval'
);

-- Later, approve from anywhere
SELECT df.signal('a1b2c3d4', 'approval',
                 '{"approved": true, "approver": "jane@contoso.com"}');

Dayanıklı HTTP çağrısı

df.http() dış çağrıları dayanıklı etkinlikler olarak yürütür; 5xx yanıtları, ağ hataları ve zaman aşımları otomatik olarak yeniden denenir.

SELECT df.start(
    df.http('https://api.example.com/users/123', 'GET') |=> 'user'
    ~> 'INSERT INTO users_cache (data) VALUES (($user::jsonb->>''body'')::jsonb)',
    'fetch-user'
);

pg_durable'de izin verilen HTTP güvenliği hakkında daha fazla bilgi edinin.

Gözlemleme ve çalıştırma

Her şey SQL'den sorgulanabilir. Öğrenebileceğiniz ayrı bir kullanıcı arabirimi veya hizmet yoktur.

-- All instances
SELECT * FROM df.list_instances();

-- Filter by status
SELECT * FROM df.list_instances() WHERE status = 'Running';
SELECT * FROM df.list_instances() WHERE status = 'Failed';

-- Detail for one instance
SELECT * FROM df.instance_info('a1b2c3d4');

-- Execution history (useful for retried or looped functions)
SELECT * FROM df.instance_executions('a1b2c3d4', 20);

-- The function graph as it ran
SELECT * FROM df.instance_nodes('a1b2c3d4');

-- System-wide metrics
SELECT * FROM df.metrics();

Çalışanın hayatta olup olmadığını denetlemek için:

SELECT epoch_id, last_seen_at, now() - last_seen_at AS time_since_last_heartbeat
FROM df._worker_epoch;

time_since_last_heartbeat 15 saniyeden küçük olması, çalışanın iyi durumda olduğu anlamına gelir. Bundan daha büyük bir değer ya da hiç satır yoksa, bu işleyicinin devre dışı olduğu veya henüz başlatılmamış olduğu anlamına gelir.

Visual Studio Code iş akışlarını izleme

Visual Studio Code için PostgreSQL uzantısı, İşlem Hatları & İş Akışları görünümünde bir İş Akışları sekmesi içerir; burada pg_durable iş akışı örneklerini inceleyebilir ve düzenleyicinizden yürütme durumunu izleyebilirsiniz.

İş Akışları bölmesini açma

  1. Visual Studio Code'da PostgreSQL uzantısını açın.
  2. Nesne Gezgini'da veritabanınıza sağ tıklayın.
  3. İşlem Hatları ve İş Akışları'ı seçin.
  4. İş Akışları sekmesini seçin.

Sol bölmede PG Dayanıklı Çalıştırmalar listelenir ve orta bölmede seçili iş akışı örneğinin ayrıntıları gösterilir.

Visual Studio Code’daki PostgreSQL uzantısının İş Akışları sekmesinin, PG Durable Runs ve iş akışı ayrıntılarını gösteren ekran görüntüsü.

İş akışı çalıştırmalarını inceleyin

Bir iş akışı çalıştırmasını seçtiğinizde, aşağıdakileri doğrulamak için özeti gözden geçirin:

  • Durum: completed, running veya failed.
  • Çalıştırma Kimliği: Örnek için benzersiz tanımlayıcı.
  • Başlangıç zamanı ve süresi: Yürütme ilerleme durumunu ve performansını izleyin.
  • Ayrıntılar paneli: Ek yürütme meta verileri.

Daha derine inmek için kullanılabilir sekmeleri kullanın:

  • Graf: İş akışı yapısını ve adım akışını gösteren görsel adım adım yürütme görünümü.
  • Zamanlama: Performans analizi ve performans sorunu belirleme için süre odaklı görünüm.
  • Sonuçlar: İş akışı yürütmesinden çıkış ve sonuç odaklı ayrıntılar.

Yapay zeka işlem hatlarıyla ilgili iş akışları için, bir İşlem hattı tanımını görüntüle eylemi (kullanılabilir olduğunda), bir iş akışı çalıştırmasından işlem hattı tanımına geri bağlanmanızı sağlar; çalıştırmalar arasındaki davranışı karşılaştırmak veya regresyonları araştırmak için kullanışlıdır.

Kimlik ve yalıtım

Dayanıklı işlevler, çalışanın ayrıcalıklarıyla değil, gönderen kullanıcının ayrıcalıklarıyla yürütülür. pg_durable, gönderim sırasında hem session_user hem de current_user öğelerini yakalar; bu nedenle SET ROLE bağlamında gönderilen işlevler bu etkin rolle çalışır.

Bu, şu anlama gelir:

  • Kullanıcılar yalnızca erişim izinleri olan verileri görür ve değiştirir.
  • Süper kullanıcı olmayanlar, kalıcı bir işlev gönderme yoluyla ayrıcalık yükseltemez.
  • Rolünüz ve verme modeliniz doğru olduğu sürece çok kiracılı iş yükleri yalıtılmış durumda kalır.

Çoğaltmalar, yedekleme ve PITR ile etkileşim

  • Yedekleme ve PITR. İşlev grafı (df şema) ve yürütme durumu (duroxide şema) normal tablolarda depolanır ve HorizonDB yedeklemelerine dahil edilir. Belirli bir zamandaki geri yükleme, her ikisini de geri yükler.
  • Okuma amaçlı çoğaltmalar. Arka plan çalışanı yalnızca birincil üzerinde çalışır. Okuma çoğaltmaları df.* izleme görünümlerini sorgulayabilir, ancak işlevleri çalıştırmaz.
  • Yük devretme. Yük devretme işleminden sonra, yeni birincildeki çalışan eski birincilin kaldığı yerden devam eder. Çalışan örnekler son denetim noktalarından yeniden başlar.

Dış düzenleyicilerle karşılaştırıldığında

Aspect Dış düzenleyici pg_durable
Deployment Ayrı hizmet, ayrı kimlik, ayrı durum deposu Bir veritabanı
Durum kalıcılığı Orchestrator'ın depolama katmanı Verileriniz için olduğu gibi aynı yedeklemeler, HA ve PITR
Identity Çalışanlar bir hizmet kimliği altında çalışıyor İşlevler gönderen kullanıcı olarak yürütülür
Hata modları Orchestrator ile veritabanı arasındaki ağ Yok - aynı işlem
En iyi kullanım alanları Birçok hizmete dokunan sistem arası düzenleme İşin büyük bölümünün Postgres içinde veya yakınında olduğu iş yükleri

pg_durable , sistemler arası işlem hatları için dış düzenleyicileri değiştirmeyi denemez. Çalışmanın çoğu veritabanı çalışması olduğunda (eklemeler, dönüşümler, yapay zeka çağrıları, zamanlanmış bakım) doğru seçimdir ve başka bir hizmet eklemek avantajdan daha maliyetlidir.

Önizleme sırasındaki sınırlamalar

  • df.http() 5xx ve ağ hatalarını yeniden deneme. 4xx yanıtları, sizin işlemeniz için iş akışına geri döndürülür; otomatik olarak yeniden denenmez.
  • Arka plan çalışanı örnek başına tek bir veritabanı sunar. Çok veritabanılı fan-out, çalışanın veritabanında çalışan bir işlev aracılığıyla df.start(..., database => 'other_db') desteklenir.
  • İşlev tanımları ve yürütme durumu, pg_durable sırasında 'ın büyük sürümleri arasında aktarılabilir değildir. Yükseltmeden önce çalışan örnekleri boşaltın veya iptal edin.