From ffa1895620db94aea3b596300869b56a01878c90 Mon Sep 17 00:00:00 2001 From: Tsiry Sandratraina Date: Sat, 17 Jan 2026 08:47:55 +0300 Subject: [PATCH] Improve batch flushing and DB pragmas Serialize batch flushes with a flushPromise and requeue failed events to enable retries. Clear timers and flush immediately when batch is full. Reduce batch size to 100 and increase timeout to 100ms. Update SQLite PRAGMAs: busy_timeout=30000, wal_autocheckpoint=1000, temp_store=MEMORY --- tap/src/drizzle.ts | 9 +++++- tap/src/tap.ts | 68 ++++++++++++++++++++++++++++------------------ 2 files changed, 49 insertions(+), 28 deletions(-) diff --git a/tap/src/drizzle.ts b/tap/src/drizzle.ts index f8a6fde8..77fea9a8 100644 --- a/tap/src/drizzle.ts +++ b/tap/src/drizzle.ts @@ -6,10 +6,17 @@ const client = createClient({ }); await client.execute("PRAGMA journal_mode = WAL;"); -await client.execute("PRAGMA busy_timeout = 5000;"); + +await client.execute("PRAGMA busy_timeout = 30000;"); + await client.execute("PRAGMA synchronous = NORMAL;"); + await client.execute("PRAGMA cache_size = -10000;"); +await client.execute("PRAGMA wal_autocheckpoint = 1000;"); + +await client.execute("PRAGMA temp_store = MEMORY;"); + const db = drizzle(client); export default { db }; diff --git a/tap/src/tap.ts b/tap/src/tap.ts index ff782205..919ffc6b 100644 --- a/tap/src/tap.ts +++ b/tap/src/tap.ts @@ -8,8 +8,8 @@ import type { InsertEvent } from "./schema/event.ts"; export const TAP_WS_URL = Deno.env.get("TAP_URL") || "http://localhost:2480"; -const BATCH_SIZE = 200; -const BATCH_TIMEOUT_MS = 50; +const BATCH_SIZE = 100; +const BATCH_TIMEOUT_MS = 100; export default function connectToTap() { const tap = new Tap(TAP_WS_URL); @@ -17,50 +17,64 @@ export default function connectToTap() { let eventBatch: InsertEvent[] = []; let batchTimer: number | null = null; - let isFlushingBatch = false; + let flushPromise: Promise | null = null; async function flushBatch() { - if (eventBatch.length === 0 || isFlushingBatch) return; - - isFlushingBatch = true; - const toInsert = [...eventBatch]; - eventBatch = []; - - try { - logger.info`🔄 Flushing batch of ${toInsert.length} events...`; - - const results = await ctx.db - .insert(schema.events) - .values(toInsert) - .onConflictDoNothing() - .returning() - .execute(); + if (flushPromise) { + await flushPromise; + return; + } - for (const result of results) { - broadcastEvent(result); + if (eventBatch.length === 0) return; + + flushPromise = (async () => { + const toInsert = [...eventBatch]; + eventBatch = []; + + try { + logger.info`🔄 Flushing batch of ${toInsert.length} events...`; + + const results = await ctx.db + .insert(schema.events) + .values(toInsert) + .onConflictDoNothing() + .returning() + .execute(); + + for (const result of results) { + broadcastEvent(result); + } + + logger.info`📝 Batch inserted ${results.length} events`; + } catch (error) { + logger.error`Failed to insert batch: ${error}`; + // Re-add failed events to the front of the batch for retry + eventBatch = [...toInsert, ...eventBatch]; + } finally { + flushPromise = null; } + })(); - logger.info`📝 Batch inserted ${results.length} events`; - } catch (error) { - logger.error`Failed to insert batch: ${error}`; - } finally { - isFlushingBatch = false; - } + await flushPromise; } function addToBatch(event: InsertEvent) { eventBatch.push(event); + // Clear existing timer if (batchTimer !== null) { clearTimeout(batchTimer); + batchTimer = null; } + // Flush immediately if batch is full if (eventBatch.length >= BATCH_SIZE) { flushBatch().catch((err) => logger.error`Flush error: ${err}`); } else { + // Set timer to flush after timeout batchTimer = setTimeout(() => { - flushBatch().catch((err) => logger.error`Flush error: ${err}`); batchTimer = null; + flushBatch().catch((err) => logger.error`Flush error: ${err}`); }, BATCH_TIMEOUT_MS); } } -- 2.51.2