BIAnalytics yükleniyor

Artımlı yüklemede kaybolan kayıtlar: updated_at'in 5 açığı

updated_at filtreli artımlı yüklemede hedefte satır eksik kalıyorsa nedeni bu beş açıktan biridir. PostgreSQL sorguları, düzeltmeler ve karar tablosu.

BIAnalytics BIAnalytics

WHERE updated_at > :son_watermark filtresi sessizce satır kaçırır ve hata vermez. Kayıp beş yerden olur: aynı zaman damgasını paylaşan satırlar, commit'i geç gelen işlemler, örtüşen pencerede çift yazma, silinen satırlar ve updated_at'i hiç güncellemeyen yazıcılar. Aşağıda her biri için bir sorgu ve bir düzeltme var. Artımlı yüklemenin ne olduğunu biliyorsunuz varsayıyorum; tanım için artımlı yükleme rehberine bakabilirsiniz.

Denemek için iki tablo ve bir watermark kaydı

Varsayalım kaynakta sipariş tablosu, hedefte onun kopyası var. Verilerin tamamı örnektir.

CREATE TABLE siparis (
  id          bigint PRIMARY KEY,
  tutar       numeric(12,2) NOT NULL,
  durum       text NOT NULL,
  updated_at  timestamptz NOT NULL DEFAULT now(),
  silindi_mi  boolean NOT NULL DEFAULT false
);
CREATE TABLE hedef_siparis (LIKE siparis INCLUDING ALL);
CREATE TABLE yukleme_durumu (
  tablo      text PRIMARY KEY,
  wm_zaman   timestamptz NOT NULL,   -- okunan son satırın updated_at değeri
  wm_id      bigint NOT NULL         -- okunan son satırın id değeri (paketli çekimde şart)
);
INSERT INTO yukleme_durumu VALUES ('siparis', '2026-09-30 00:00:00+03', 0);

Aynı zaman damgalı satırlardan biri kaybolur

Tek bir işlemde yazılan satırların hepsi aynı now() değerini taşır. Yükleme toplu çekiyorsa (LIMIT) ve watermark'ı son satırın updated_at değeri yapıyorsa, sınırda kalan kardeş satırlar bir sonraki turda > filtresine takılır.

INSERT INTO siparis (id, tutar, durum, updated_at) VALUES
  (1, 100, 'yeni', '2026-10-01 10:00:00+03'),
  (2, 250, 'yeni', '2026-10-01 10:00:00+03'),
  (3,  75, 'yeni', '2026-10-01 10:00:00+03');

-- 1. tur, 2'şerli paket: id 1 ve 2 gelir, watermark 10:00:00 olur
SELECT id FROM siparis
WHERE updated_at > '2026-09-30 00:00:00+03'
ORDER BY updated_at, id LIMIT 2;

-- 2. tur: 0 satır döner, id 3 hedefe hiç gitmez
SELECT id FROM siparis WHERE updated_at > '2026-10-01 10:00:00+03';

İlk akla gelen düzeltme sınırı >= yapmaktır, ama paketli çekimde işe yaramaz: aynı örnekte 2. turu >= ile çalıştırdım, yine id 1 ve 2 döndü, watermark 10:00:00'da kaldı ve id 3'e sıra gelmedi. Yükleme sonsuza dek aynı paketi okur.

>= yalnız çekim LIMIT'siz ve hedef yazma ON CONFLICT ile idempotent ise doğru çözümdür: sınır satırı iki kez gelir, id anahtarıyla üzerine yazılır; geç commit payı için pencere bölümüne bakın. Sütun hassasiyeti düşükse (saniye gibi) paket sınırı olmasa da benzer kayıp çıkabilir: tur bittikten sonra commit edilen, aynı saniyeyi taşıyan satır > filtresine takılır; bu durum da aynı gruba girer.

Tercihim net: yüklemeniz LIMIT ile paketliyorsa >= ve örtüşen pencere yerine bileşik anahtar ve üst sınır kullanın (pencere paketle ilerlemeyi durdurur; pencere bölümünde anlatılıyor). Sıralamayı ve filtreyi (updated_at, id) ikilisiyle kurun, watermark'ı son satırın iki değerine birden güncelleyin. Üst sınır, açık işlemdeki satırı beklettiği için sonraki bölümdeki geç commit kaybını önler; pay örnek olarak 15 dakikadır ve o bölümdeki ölçümle seçilir. Bedeli: yükleme hedefe pay kadar geriden gelir.

-- 1. tur: watermark (2026-09-30, 0) → id 1 ve 2 gelir
SELECT s.id, s.updated_at
FROM siparis s
JOIN yukleme_durumu y ON y.tablo = 'siparis'
WHERE (s.updated_at, s.id) > (y.wm_zaman, y.wm_id)
  AND s.updated_at < now() - interval '15 minutes'   -- pay: örnek değer
ORDER BY s.updated_at, s.id
LIMIT 2;

-- paket yazılınca watermark = son satır (10:00:00, 2); 2. turda id 3 gelir
UPDATE yukleme_durumu SET wm_zaman = '2026-10-01 10:00:00+03', wm_id = 2
WHERE tablo = 'siparis';

PostgreSQL 14'te üç turluk tam döngüyü üst sınırsız sorguyla çalıştırdım: paketler 2, 1, 0 satır çıktı, hedefe id 1, 2 ve 3 ulaştı. Üst sınırlı sorguyu tam döngüde değil, yalnız sonraki bölümdeki iki oturumlu senaryoda çalıştırdım. Paketi okuma, yazma ve watermark güncellemesini tek işlemde çalıştırın.

Commit'i geç gelen satır pencerenin altında kalır

now() işlemin başladığı anı döndürür; işlem boyunca değişmez. PostgreSQL belgesi bunu işlemin başlangıç zamanı olarak tanımlar, clock_timestamp() ise gerçek anı verir. İkisi de commit anı değildir. Read Committed düzeyinde bir sorgu ise yalnız başlamadan önce commit edilmiş veriyi görür: açık işlemdeki satırı görmez.

İki psql oturumuyla görülür. Önceki turun watermark'ı 09:59:00 olsun; id 10'un damgası bu sınırın üstündedir, yani tek engel işlemin henüz commit edilmemiş olmasıdır. Saatler örnektir.

-- Oturum A (10:00:00)
BEGIN;
INSERT INTO siparis (id, tutar, durum) VALUES (10, 500, 'yeni');
-- updated_at = 10:00:00; COMMIT henüz yok, başka iş sürüyor

-- Oturum B (10:00:20)
INSERT INTO siparis (id, tutar, durum) VALUES (11, 80, 'yeni');  -- updated_at = 10:00:20

-- Yükleme (10:00:30): yalnız id 11'i görür, watermark = 10:00:20
SELECT id, updated_at FROM siparis WHERE updated_at > '2026-10-01 09:59:00+03';

-- Oturum A (10:01:00)
COMMIT;   -- id 10 artık görünür ama updated_at = 10:00:00 < watermark

İki gerçek oturumla çalıştırdım: yükleme, bu örnekte eklenen satırlardan yalnız id 11'i döndürdü ve watermark 10:00:20 oldu. A commit ettikten sonra yeni watermark'la çalışan filtre (updated_at > 10:00:20) 0 satır döndürdü; id 10 sonsuza dek dışarıda kalır. Düzeltme:

  1. Watermark'ı "turun bittiği an" değil, "okunan satırlardaki en büyük updated_at" olarak saklayın.
  2. Bir güvenlik payı koyun. Paketsiz çekimde tur watermark'tan pay kadar geriden başlar: updated_at >= watermark - interval '15 minutes' (pencere bölümü). Paketli çekimde pencere kullanılmaz; yeni satırlar bekletilir: updated_at < now() - interval '15 minutes' (önceki bölüm). İkisi aynı sorguda birlikte kullanılmaz: pencere geriye bakar, üst sınır yeni satırı bekletir.
  3. Payı, kaynaktaki en uzun işlemden büyük seçin; 15 dakika yalnız örnektir. Bedeli: üst sınırlı yükleme hedefe pay kadar geriden gelir. Fikir için SELECT max(now() - xact_start) FROM pg_stat_activity WHERE xact_start IS NOT NULL AND pid <> pg_backend_pid(); anlık bir ölçüm verir (xact_start işlemin başlangıcıdır; pid koşulu kendi oturumunuzu saymaz). Gece toplu işleri varsa gün boyu tekrar bakın.

Pencereyi örtüştürünce çift yazmamak

Örtüşen pencere satırları tekrar getirir; hedefin buna dayanması gerekir. Burada iki tuzak var. Birincisi, aynı INSERT komutunda aynı anahtardan iki satır gelirse hata alırsınız. Kaynak bir değişim günlüğü ya da join sonucuysa bu olur.

INSERT INTO hedef_siparis (id, tutar, durum, updated_at) VALUES
  (1, 100, 'yeni',  '2026-10-01 10:00:00+03'),
  (1, 100, 'odendi','2026-10-01 10:05:00+03')
ON CONFLICT (id) DO UPDATE SET durum = EXCLUDED.durum;

PostgreSQL'in verdiği hata ON CONFLICT DO UPDATE command cannot affect row a second time (SQLSTATE 21000). Belge DO UPDATE'i deterministik bir komut sayar: komut aynı satırı birden çok kez etkileyemez, aksi hâlde cardinality violation hatası çıkar. İkinci tuzak, geç gelen eski sürümün hedefteki yeni sürümü ezmesidir. İkisini birlikte çözen kalıp:

INSERT INTO hedef_siparis (id, tutar, durum, updated_at, silindi_mi)
SELECT DISTINCT ON (id) id, tutar, durum, updated_at, silindi_mi
FROM siparis
WHERE updated_at >= (SELECT wm_zaman FROM yukleme_durumu WHERE tablo = 'siparis')
                    - interval '15 minutes'
ORDER BY id, updated_at DESC
ON CONFLICT (id) DO UPDATE
SET tutar = EXCLUDED.tutar, durum = EXCLUDED.durum,
    updated_at = EXCLUDED.updated_at, silindi_mi = EXCLUDED.silindi_mi
WHERE hedef_siparis.updated_at <= EXCLUDED.updated_at;

DISTINCT ON (id) ve ORDER BY id, updated_at DESC anahtar başına en yeni sürümü bırakır. WHERE, belgenin dediği gibi çakışma bulunduktan sonra değerlendirilir ve yalnız koşulu sağlayan satırlar güncellenir. <= kullanıldı çünkü aynı zaman damgalı ama farklı içerikli sürüm de uygulanmalı; bunun bedeli, değişmemiş satırı yeniden yazmaktır.

Tur bitince watermark'ı, yüklemenin az önce okuduğu satırlardaki en büyük updated_at değerine çekin; geri gitmesin diye greatest kullanın. Payı watermark'a değil, bir sonraki turun filtresine uygulayın. Değeri kaynağa ikinci bir max() sorgusuyla almayın: iki sorgu arasında yeni satır commit olabilir ve watermark okunmamış satırın önüne geçer.

UPDATE yukleme_durumu
SET wm_zaman = greatest(wm_zaman, :'okunan_en_buyuk_updated_at')
WHERE tablo = 'siparis';

Bu kalıp LIMIT ile paketlenmez: pencere her turda pay kadar geriye baktığı için aynı satırlar paketin başına döner ve ilerleme durabilir. Paketlemeniz gerekiyorsa pencere yerine bileşik anahtar ve üst sınır (updated_at < now() - pay) kullanın; pay, yukarıdaki ölçümle seçilir.

Silinen satır sorguda hiç görünmez

DELETE ile giden satır, updated_at filtresinin bakacağı bir iz bırakmaz. Üç yol var:

  • Yumuşak silme: uygulama DELETE yerine silindi_mi = true ve updated_at güncellemesi yapar. Yukarıdaki yükleme bunu zaten taşır.
  • Anahtar mutabakatı: kaynaktaki anahtar listesi ile hedeftekini karşılaştırın; hedefte olup kaynakta olmayanı silin ya da işaretleyin.
SELECT h.id
FROM hedef_siparis h
LEFT JOIN siparis s ON s.id = h.id
WHERE s.id IS NULL AND NOT h.silindi_mi;

CDC için Debezium seçerseniz iki ayrıntı hazırlıksız yakalar. Debezium bir silme için iki olay üretir: "op": "d" olayı ve aynı anahtarlı, değeri null olan bir tombstone olayı; tombstones.on.delete varsayılanı true'dur. Tüketici ikisini ayrı ele almalıdır. Ayrıca PostgreSQL bağlayıcısı belgesine göre REPLICA IDENTITY DEFAULT iken silme olayında önceki değer olarak yalnız birincil anahtar gelir; güncellemede önceki değer yalnız anahtar değiştiyse gelir (PostgreSQL mantıksal çoğaltma mesaj biçimi, Update mesajı). FULL ise tüm sütunların önceki değerini taşır. Silinen satırın diğer sütunlarını da görmek için ALTER TABLE siparis REPLICA IDENTITY FULL; çalıştırın.

updated_at'i uygulama güncellemiyorsa

Bir betik, elle yapılan düzeltme ya da başka bir servis UPDATE çalıştırıp sütuna dokunmazsa satır filtreye hiç girmez. Çözüm sütunu veritabanına yazdırmaktır:

CREATE FUNCTION updated_at_yaz() RETURNS trigger
LANGUAGE plpgsql AS $$
BEGIN
  NEW.updated_at := clock_timestamp();
  RETURN NEW;
END $$;

CREATE TRIGGER siparis_updated_at
BEFORE INSERT OR UPDATE ON siparis
FOR EACH ROW EXECUTE FUNCTION updated_at_yaz();

BEFORE tetikleyicisi eklenecek ya da güncellenecek satırı değiştirebilir. clock_timestamp() seçildi, çünkü aynı işlemdeki satırlar farklı değerler alır ve ilk açığın paket sınırı sorunu azalır. Ama tetikleyici ikinci açığı kapatmaz: damga yine commit anı değildir. Pencere ya da üst sınır yerinde kalır.

Zaman damgası mı CDC mi: dört soruda karar

Soruupdated_at + örtüşen pencereCDC
Kaynakta satır siliniyor mu?Yalnız yumuşak silme ya da mutabakat sorgusuylaSilmeyi doğrudan taşır
Kabul edilen gecikmeTur aralığı kadar. Paketsiz pencere yüklemeyi geciktirmez; üst sınırlı paketli çekim hedefe pay kadar geriden gelirDeğişiklik günlükten akar; tur aralığı beklenmez
Kaynakta mantıksal çoğaltma açılabilir mi?Gerekmezwal_level = logical ve max_replication_slots en az 1 olmalı (belge)
Tablo büyüklüğü ve yazma hızıupdated_at üzerinde indeks yoksa sorgu planı genellikle tabloyu baştan sona okur; önce indeks eklenip eklenemeyeceğine bakınDeğişiklikler kaynağın yazma günlüğünden (WAL) okunur; yük tablo taramasından günlük okumasına ve yuvaya kayar

CDC'nin bedeli kaynak tarafındadır. PostgreSQL belgesi yuvanın tüketicisi durmuş olsa bile WAL'ı tutmaya devam ettiğini yazar; kullanılmayan yuva diski doldurabilir. Yuvayı izleyecek kişi yoksa CDC'ye geçmek riski yer değiştirir.

Üç tercih:

  • Silme yok ve updated_at güvenilirse: zaman damgası ve ON CONFLICT yeter; paketsizse örtüşen pencere, paketliyse bileşik anahtar ve üst sınır ekleyin. Pay, geç commit bölümündeki ölçümle seçilir (15 dakika örnektir); üst sınır hedefi pay kadar geriletir.
  • Silme var, kaynakta wal_level değiştirilebilir ve yuva izlenebilir: CDC.
  • Silme var ama kaynağa dokunamıyorsanız: yumuşak silme ya da günlük mutabakat; ikisi de olmuyorsa haftalık tam yükleme planlayın.

Yükleme, doğrulama ve hata yönetimi adımlarını dışarıdan değerlendirtmek isterseniz ETL danışmanlığı sayfasında kapsam anlatılıyor.

Yükleme bitti, hedefte satır eksik mi

Yüklemenin son adımı gün bazlı bir sayım karşılaştırması olsun. Eşik aşılırsa hedefi "yayımlandı" işaretlemeyin.

WITH k AS (
  SELECT (updated_at AT TIME ZONE 'UTC')::date AS gun,
         count(*) AS adet, max(updated_at) AS son
  FROM siparis GROUP BY 1
), h AS (
  SELECT (updated_at AT TIME ZONE 'UTC')::date AS gun,
         count(*) AS adet, max(updated_at) AS son
  FROM hedef_siparis GROUP BY 1
)
SELECT coalesce(k.gun, h.gun) AS gun,
       k.adet AS kaynak_adet, h.adet AS hedef_adet,
       k.son AS kaynak_son, h.son AS hedef_son
FROM k FULL JOIN h ON h.gun = k.gun
WHERE k.adet IS DISTINCT FROM h.adet
   OR k.son  IS DISTINCT FROM h.son
ORDER BY 1;

Gün hesabında AT TIME ZONE 'UTC' var: ::date oturumun saat dilimine bağlıdır ve iki bağlantı farklı ayarlıysa gün sınırındaki satırlar farklı güne düşer. Sorgu boş dönerse iki taraf bu iki ölçüde uyumludur; satır dönerse o günü örtüşen pencereyle yeniden yükleyin. Üst sınırlı yüklemede hedef pay kadar geridedir; sorguya da aynı sınırı koymazsanız bugünün son dakikaları yanlış alarm verir.

Bu yazıyı paylaşın LinkedIn'de paylaş

Burada anlatılana benzer bir işiniz mi var?

Yazıların arkasındaki ekibe doğrudan yazabilirsiniz.

Bizimle iletişime geçin

Demo Talebi

Çözümlerimizi yakından tanımak için formu doldurun.

person
mail
phone
notes