Pelacak ekspedisi laut jarak jauh mengandalkan uplink satelit berbiaya tinggi dan berlatensi tinggi. Saat jaringan terputus di tengah laut, perangkat menyimpan log telemetri di memori lokal. Ketika koneksi pulih, mekanisme retransmisi otomatis membanjiri Ingestion API dengan paket-paket tertunda.
Masalah timbul saat paket yang dikirim belakangan tiba lebih cepat dibanding paket sebelumnya akibat variasi rute jaringan atau konkurensi retry. Jika API langsung menimpa data perangkat dengan payload yang baru masuk, terjadi state drift: koordinat lama menimpa koordinat terkini di basis data operasional.
1. Monotonic Sequence Number vs Client Timestamp
Mengandalkan timestamp dari perangkat untuk mengurutkan state adalah anti-pattern pada sistem telemetri remote. Masalah utamanya meliputi:
- Clock drift: Real-Time Clock (RTC) perangkat berbiaya rendah dapat bergeser beberapa detik hingga menit per bulan.
- GPS Lock Latency: Waktu sistem firmware sering kali tidak sinkron sampai cold-start GPS memperoleh lock. Selama transisi ini, log awal memakai timestamp default sistem (misal 1970-01-01).
- Manipulasi Waktu: Kegagalan sinkronisasi NTP akibat keterbatasan bandwidth satelit.
Solusi deterministik: gunakan Monotonic Sequence Number (integer 64-bit). Firmware menginkremen angka ini secara linier untuk setiap record metrik yang dibuat. Backend menentukan urutan keabsahan state murni berdasarkan nilai sekuens, bukan timestamp.
// Kontrak Payload Ingestion
{
"device_id": "track-exp-084",
"batch_id": "b9f2c8d4-53c1-4828-98e6-71d3df85223a",
"records": [
{
"seq": 10420,
"recorded_at": "2026-03-31T04:12:00Z",
"lat": 21.3069,
"lon": -157.8583,
"battery": 3.82
}
]
}2. Idempotency Key per Batch
Saat ACK HTTP gagal mencapai perangkat akibat koneksi terputus tiba-tiba, firmware akan mengirim ulang batch yang sama. Untuk mencegah pemrosesan ganda, sertakan batch_id (UUID v4 yang dibuat sebelum pengiriman pertama).
Simpan batch_id di cache terdistribusi (Redis) atau basis data transaksional dengan TTL terukur:
- Jika
batch_idsudah ada: API segera mengembalikan status200 OKtanpa memproses ulang komputasi state. - Jika
batch_idbaru: simpan ID secara atomik, lalu proses payload ke storage engine.
3. Conditional Write (Optimistic Concurrency Control)
Penyimpanan telemetri biasanya terbagi dua: time-series/append-only log (menyimpan seluruh jejak historis) dan materialized current state (menyimpan posisi terbaru saat ini untuk dashboard pemantau). Masalah drift hanya terjadi pada tabel current state.
Gunakan operasi Compare-And-Swap (CAS) langsung di persistence layer menggunakan klausa SQL kondisional:
UPDATE device_current_state
SET
last_seq = :seq,
latitude = :lat,
longitude = :lon,
battery = :battery,
updated_at = NOW()
WHERE
device_id = :device_id
AND last_seq < :seq;Evaluasi hasil eksekusi kueri:
- Affected rows = 1: Paket valid dan merupakan data terbaru. State berhasil diperbarui.
- Affected rows = 0: Paket usang (out-of-order) atau duplikat. Penting: Tetap kembalikan respons
200 OKatau202 Acceptedke perangkat. Jika API mengembalikan error4xxatau5xx, firmware pelacak akan terus mencoba mengirim ulang paket kadaluwarsa tersebut, menghabiskan kuota data dan daya baterai.
4. Implementasi Handler dan runnable Assert Test
Berikut implementasi backend minimal menggunakan Python standard library (SQLite in-memory) yang memvalidasi deduplikasi batch dan penolakan payload out-of-order secara aman.
import sqlite3
from typing import Dict, Any, Tuple
class TelemetryIngestor:
def __init__(self, db_conn: sqlite3.Connection):
self.db = db_conn
self._init_db()
def _init_db(self):
with self.db:
self.db.execute("""
CREATE TABLE IF NOT EXISTS processed_batches (
batch_id TEXT PRIMARY KEY
);
""")
self.db.execute("""
CREATE TABLE IF NOT EXISTS device_current_state (
device_id TEXT PRIMARY KEY,
last_seq INTEGER NOT NULL,
lat REAL NOT NULL,
lon REAL NOT NULL
);
""")
def ingest(self, payload: Dict[str, Any]) -> Tuple[int, str]:
device_id = payload["device_id"]
batch_id = payload["batch_id"]
records = payload["records"]
with self.db:
# 1. Cek & Simpan Idempotency Key
cursor = self.db.execute(
"INSERT OR IGNORE INTO processed_batches (batch_id) VALUES (?)",
(batch_id,)
)
if cursor.rowcount == 0:
# Batch sudah pernah diproses, ACK langsung
return 200, "DUPLICATE_BATCH_ACK"
# 2. Proses record menggunakan Conditional Update (CAS)
for rec in records:
seq = rec["seq"]
lat = rec["lat"]
lon = rec["lon"]
# Upsert jika device baru; Update jika seq lebih besar
# ponytail: Gunakan CTE/ON CONFLICT native PostgreSQL untuk skenario produksi
self.db.execute(
"INSERT OR IGNORE INTO device_current_state (device_id, last_seq, lat, lon) VALUES (?, ?, ?, ?)",
(device_id, seq, lat, lon)
)
self.db.execute(
"""
UPDATE device_current_state
SET last_seq = ?, lat = ?, lon = ?
WHERE device_id = ? AND last_seq < ?
""",
(seq, lat, lon, device_id, seq)
)
return 200, "PROCESSED"
# ==================== RUNNABLE SELF-CHECK TEST ====================
if __name__ == "__main__":
conn = sqlite3.connect(":memory:")
ingestor = TelemetryIngestor(conn)
device = "hawaii-tracker-01"
# Paket 1: Normal (Seq 100)
p1 = {
"device_id": device,
"batch_id": "b-001",
"records": [{"seq": 100, "lat": 21.30, "lon": -157.85}]
}
# Paket 2: Terlambat / Re-ordered retry dari masa lalu (Seq 90)
p2_late = {
"device_id": device,
"batch_id": "b-000-late",
"records": [{"seq": 90, "lat": 20.50, "lon": -156.00}]
}
# Paket 3: Update terbaru (Seq 105)
p3_new = {
"device_id": device,
"batch_id": "b-002",
"records": [{"seq": 105, "lat": 22.00, "lon": -159.00}]
}
# Eksekusi p1
status, msg = ingestor.ingest(p1)
assert status == 200 and msg == "PROCESSED"
# Eksekusi p2_late (harus di-drop secara silent di tabel state)
status, msg = ingestor.ingest(p2_late)
assert status == 200 and msg == "PROCESSED"
# Cek State setelah p2: Posisi HARUS tetap Seq 100, bukan 90
cur = conn.cursor()
cur.execute("SELECT last_seq, lat FROM device_current_state WHERE device_id = ?", (device,))
row = cur.fetchone()
assert row[0] == 100, f"State drift terdeteksi! seq={row[0]}"
assert row[1] == 21.30, f"Koordinat tertimpa! lat={row[1]}"
# Eksekusi p3_new (harus mengupdate state)
ingestor.ingest(p3_new)
cur.execute("SELECT last_seq, lat FROM device_current_state WHERE device_id = ?", (device,))
row = cur.fetchone()
assert row[0] == 105
assert row[1] == 22.00
# Eksekusi ulang p3_new (Idempotency check)
status, msg = ingestor.ingest(p3_new)
assert status == 200 and msg == "DUPLICATE_BATCH_ACK"
print("Semua assert lolos: State drift berhasil dicegah.")
Trade-offs dan Pertimbangan Sistem
- Historical vs Current State: Conditional write hanya diterapkan pada tabel representasi state terkini. Tabel append-only / historical log (TimescaleDB, ClickHouse) tetap menerima seluruh payload agar trek rute lengkap tidak hilang.
- Sequence Reset: Jika firmware di-flash ulang dan sequence kembali ke
0, backend memerlukan endpoint administrasi untuk meresetlast_seq, atau menyertakansession_id/epoch_idunik bersamaan dengan sequence number.
Komentar
0 komentar
Masuk ke akun kamu untuk ikut berkomentar.
Belum ada komentar
Jadilah yang pertama ikut berdiskusi!