Tarih Bölümlü Veri Hatlarında Yeniden Başlatmaya Dayanıklı Backfill Tasarımı
Tarih bölümlü veri hatlarında backfill işleminin güvenilirliği, ilerleme konumunu bellekte tutmaktan çok commit edilmiş çıktılardan durum türetmeye dayanır. Bu tasarım; veri kapanış sınırı, ters kronolojik eksik bölüm taraması, atomik yayın ve idempotent etki ilkelerini birlikte ele alır.
[d, d + 1) aralığındaki bir günün işlenebilir olması, yalnızca takvim gününün bitmesine değil, kaynak verinin kapanış politikasına bağlıdır. Günlük gecikme toleransı L, mevcut zaman t ve gün sonu end(d) için koşul şöyledir:
eligible(d, t) = t >= end(d) + L
Buna göre en son kapanmış gün:
Dclosed(t) = max { d | eligible(d, t) }
Bir günlük çıktı dosyasının varlığı ancak üretimin eksiksiz tamamlandığını kesin olarak temsil ediyorsa güvenilir durum bilgisidir. Hedef dosya yazım sırasında görünür hâle geliyorsa, süreç yarıda kesildiğinde kısmi içerik bırakıyorsa veya ilerleme yalnızca bellekte tutuluyorsa, yeniden başlatma sonrasında hangi günlerin gerçekten tamamlandığı belirlenemez.
Veri kapanış sınırı ve işlem zamanı
Takvim gününün sona ermesi, o güne ait verinin işlenmeye hazır olduğu anlamına gelmez. Kaynak sistemde işlemler gecikmeli tamamlanabilir, son kayıtlar gece yarısından sonra gelebilir veya günlük veri belirli bir operasyonel saatte kararlı hâle gelebilir.
Örneğin kesinleştirme sınırı ertesi gün saat 01:00 ise, 4 Ağustos saat 00:40 itibarıyla 3 Ağustos henüz kapalı kabul edilmez; en son güvenli gün 2 Ağustos'tur. Saat 01:00 sonrasında 3 Ağustos işlenebilir hâle gelir.
Bu sınır, stream processing sistemlerindeki watermark kavramına benzer bir işlev görür. Apache Beam'de watermark, belirli bir zaman penceresi açısından giriş verisinin ne ölçüde tamamlandığını ifade eden ilerleme göstergesidir. Günlük batch görevindeki kapanış saati ise, kaynak sistemin çalışma sözleşmesi izin veriyorsa deterministik bir zaman sınırı olabilir. Kapanıştan sonra geç veri gelebiliyorsa, sınır mutlak tamamlığı değil, iş politikasınca kabul edilen tamamlığı temsil eder.
İşlem zamanı ile veri zamanı ayrılmalıdır. Görevin 4 Ağustos'ta çalışması, mutlaka 4 Ağustos verisini işlemesi gerektiği anlamına gelmez. İşlenecek tarih, görevin çalıştığı takvim gününden değil, veri kapanış kuralından türetilir.
Durumu commit edilmiş çıktılardan türetmek
Son işlenen tarihi bir cursor ile tutmak ilk bakışta yeterli görünür:
cursorDate = 2026-07-31Ancak cursor ile gerçek çıktı durumu ayrışabilir:
- Çıktı yazılır, cursor güncellenmeden süreç kapanır.
- Cursor güncellenir, çıktı kalıcı hâle gelmeden süreç kapanır.
- Yeni bir kategori eklenir; cursor ileri tarihtedir fakat yeni kategoriye ait eski çıktılar yoktur.
- Bir dosya dışarıdan silinir; cursor bu eksikliği göstermez.
- Yapılandırma değişir; eski cursor yeni çıktı uzayını temsil etmez.
Asıl durum, hangi tarihte kalındığı değil, beklenen çıktı kümesinin hangi elemanlarının commit edilmiş olduğudur. Bir çıktı bölümü şu koordinatlarla tanımlanabilir:
P = (date, unit, categoryType, valueType)
Belirli bir tarihte beklenen bölüm kümesi:
Expected(date) = Units(date) × Queries(date)
Burada Queries, kategori ve değer türü gibi mantıksal çıktı boyutlarını kapsar. Bir bölümün tamamlanmış olması ise şu koşulla belirlenir:
complete(P) = committedArtifactExists(path(P))
Süreç her başlangıçta güncel yapılandırmadan beklenen bölümleri yeniden üretir, bunları dosya sistemindeki commit edilmiş çıktılarla karşılaştırır ve yalnızca eksik bölümleri işler. Böylece yeni bir kategori eklendiğinde Expected(date) kümesi genişler; eski tarihlerde yeni koordinatların çıktısı eksik görüldüğünden backfill kendiliğinden devreye girer.
Her sorgunun en erken anlamlı tarihi E(q) ile gösterilirse işlenecek aralık:
[max(globalStart, E(q)), Dclosed]
olur. Bu yaklaşım, sonradan eklenmiş fakat geçmiş kaynak verisi bulunmayan kategoriler için anlamsız tarihlerin taranmasını önler.
En yeni günü önceleyen deterministik tarama
Eski tarihten yeni tarihe ilerleyen klasik backfill, büyük bir tarihsel eksik olduğunda güncel veriyi geciktirebilir. Operasyonel gereksinim hem en son kapanmış günün kısa sürede yayımlanması hem de tarihsel eksiklerin zaman içinde kapatılması olduğunda, tarama ters kronolojik sırada yapılabilir.
- En son kapanmış günü belirleyin.
- Bu güne ait tüm beklenen bölümleri denetleyin.
- Eksik bölümleri tamamlayarak en güncel çıktıyı garanti edin.
- Ardından önceki güne geçin.
- En eski izin verilen tarihe kadar geriye ilerleyin.
- Kaynak, süre veya iş yükü sınırına ulaşıldığında döngüyü sonlandırın.
- Sonraki çalıştırmada aynı hesaplamayı yeniden yapın.
Temel akış şöyledir:
closedDate = latestClosedDate(now)
expected = loadCurrentDefinitions()
for date = closedDate downto earliestRelevantDate:
for partition in deterministicOrder(expected, date):
target = path(date, partition)
if committed(target):
continue
data = query(date, partition)
publishAtomically(target, data)Bu algoritmada tarih veya kategori cursor'ı güncellenmez. Her yeniden başlatmada tarama güncel kapanmış tarihten başlar. Önceden tamamlanan bölümler ucuz bir varlık denetimiyle atlanır; yarım kalan ilk bölüm yeniden üretilir.
Tarama; gün sayısı, bölüm sayısı, toplam sorgu süresi veya bir çevrimde üretilecek dosya miktarı bakımından sınırlandırılabilir. Bu sınır doğruluğu değiştirmez; yalnızca işi çevrimlere böler. Durum çıktılardan türetildiği için bu bölme ek checkpoint gerektirmez.
Atomik yayın ve commit göstergesi
Dosya varlığının commit göstergesi sayılabilmesi için doğrudan hedef dosyaya yazmak yerine geçici dosya üzerinden yayın yapılmalıdır:
veriyi üret
geçici dosyaya yaz
akışı kapat
geçici dosyayı hedef ada atomik olarak taşıJava'daki temel desen şöyledir:
final Path temporary = target.resolveSibling(target.getFileName() + ".tmp");
writeResult(temporary, result);
Files.move(temporary, target, StandardCopyOption.ATOMIC_MOVE);ATOMIC_MOVE, dosyanın atomik bir dosya sistemi işlemi olarak taşınmasını talep eder. Java API'si bu işlem desteklenmiyorsa AtomicMoveNotSupportedException üretir. POSIX rename işlemi de desteklenen koşullarda ad değişikliğinin atomik olmasını gerektirir. Hedef dosya böylece ya hiç görünmez ya da tamamlanmış hâliyle görünür; okuyucu ara yazım durumunu gözlemlemez.
Geçici dosya hedefle aynı dosya sisteminde, tercihen aynı dizinde oluşturulmalıdır. Dosya sistemleri arası taşıma, kopyalama ve silme işlemine dönüşebilir. Atomik taşımanın desteklenmediği ortamda sessizce atomik olmayan harekete geçmek, dosya varlığını commit göstergesi olarak kullanan algoritmanın varsayımını bozar.
Bu durumda iki güvenli seçenek vardır:
- Atomik yayın desteklenmiyorsa işlemi başarısız kabul etmek.
- Veri dosyasından ayrı ve yalnızca tamamlanmış yazımdan sonra oluşturulan bir commit işareti kullanmak.
İkinci modelde data.json tek başına yeterli değildir; bölüm ancak data.json.commit de mevcutsa tamamlanmış sayılır. Commit işaretinin kendisi de güvenli biçimde yayımlanmalıdır.
Bu protokolde çökme durumları iki kararlı sonuç üretir:
- Geçici dosya oluşturulmadan veya yazılırken çökülürse hedef yoktur; bölüm yeniden işlenir.
- Yazım tamamlandıktan fakat taşıma öncesinde çökülürse geçici dosya temizlenir veya üzerine yeniden yazılır.
- Atomik taşıma sonrasında çökülürse hedef vardır; bölüm tamamlanmış sayılır.
- Sonraki bölüme geçerken çökülürse tamamlanan bölüm atlanır, eksik bölümden devam edilir.
Yeniden çalıştırılabilirlik ve idempotent etki
Yeniden başlatılan bir görev aynı veritabanı sorgusunu birden fazla kez çalıştırabilir. Bu tek başına hata değildir; önemli olan kalıcı ve dışarıdan gözlemlenebilir etkinin bir kez oluşmasıdır.
Apache Flink belgelerindeki ayrımla uyumlu biçimde, exactly-once işleme her kaydın fiziksel olarak yalnızca bir kez işlenmesi anlamına gelmez. Hata sonrasında kaynak akışı geri sarılıp yeniden oynatılabilir; garanti, her olayın yönetilen durumu tam olarak bir kez etkilemesidir. Uçtan uca exactly-once için kaynağın yeniden oynatılabilir, hedefin ise transaction destekli veya idempotent olması gerekir.
Tarih bölümlü dosya üretiminde de şu ayrım geçerlidir:
exactly-once processing ≠ exactly-once observable effect
Sorgu iki kez çalışabilir; ancak deterministik çıktı aynı hedef koordinatına atomik olarak yalnızca bir commit edilmiş dosya üretir. Bu nedenle uygun niteleme, yeniden çalıştırılabilir ve idempotent yayındır.
Dosya üretimine ek olarak veritabanına kayıt yazılıyorsa, dosya sistemi taşıması ile veritabanı transaction'ı tek bir yerel transaction altında atomik değildir. Dosyayı önce commit etmek veya veritabanını önce commit etmek, her iki sırada da ara hata penceresi bırakır. Hedeflerden biri diğerinin yeniden denenmesini güvenle tolere etmelidir. Veritabanında doğal ya da yapay bir bölüm anahtarıyla idempotent insert; dosya tarafında atomik yayın ve yeniden uzlaştırma mekanizması kullanılabilir.
Durum = commit edilmiş çıktılar
time sınırı = veri kapanış politikası
ilerleme = deterministik eksik bölüm taraması
dayanıklılık = idempotent üretim + atomik yayın
Bellekte tutulan cursor bir performans optimizasyonu olabilir; ancak doğruluğun tek kaynağı olmamalıdır. Yapılandırma beklenen çıktı kümesini, dosya sistemi gerçekleşmiş çıktı kümesini ve kapanış kuralı işlenmesine izin verilen zaman aralığını tanımlar. Sistem bu kümeler arasındaki farkı işleyerek ilerler. Böyle bir görev, özel bir kurtarma senaryosuna ihtiyaç duymadan normal algoritmasını yeniden çalıştırabildiğinde hata toleransını işlem modelinin doğal parçası hâline getirir.